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 93118ed9f..73850f670 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 @@ -82,7 +82,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } if (bean instanceof ThreadPoolTaskExecutor) { if (isProxyNeeded(beanName)) { - return wrapThreadPoolTaskExecutor(bean); + return wrapThreadPoolTaskExecutor(bean, beanName); } else { log.info("Not instrumenting bean " + beanName); @@ -90,7 +90,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } else if (bean instanceof ScheduledExecutorService) { if (isProxyNeeded(beanName)) { - return wrapScheduledExecutorService(bean); + return wrapScheduledExecutorService(bean, beanName); } else { log.info("Not instrumenting bean " + beanName); @@ -98,7 +98,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } else if (bean instanceof ExecutorService) { if (isProxyNeeded(beanName)) { - return wrapExecutorService(bean); + return wrapExecutorService(bean, beanName); } else { log.info("Not instrumenting bean " + beanName); @@ -106,26 +106,26 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } else if (bean instanceof AsyncTaskExecutor) { if (isProxyNeeded(beanName)) { - return wrapAsyncTaskExecutor(bean); + return wrapAsyncTaskExecutor(bean, beanName); } else { log.info("Not instrumenting bean " + beanName); } } else if (bean instanceof Executor) { - return wrapExecutor(bean); + return wrapExecutor(bean, beanName); } return bean; } - private Object wrapExecutor(Object bean) { + private Object wrapExecutor(Object bean, String beanName) { Executor executor = (Executor) bean; boolean methodFinal = anyFinalMethods(executor); boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers()); boolean cglibProxy = !methodFinal && !classFinal; try { - return createProxy(bean, cglibProxy, - new ExecutorMethodInterceptor<>(executor, this.beanFactory)); + return createProxy(bean, cglibProxy, new ExecutorMethodInterceptor<>(executor, + this.beanFactory, beanName)); } catch (AopConfigException ex) { if (cglibProxy) { @@ -134,43 +134,43 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { "Exception occurred while trying to create a proxy, falling back to JDK proxy", ex); } - return createProxy(bean, false, - new ExecutorMethodInterceptor<>(executor, this.beanFactory)); + return createProxy(bean, false, new ExecutorMethodInterceptor<>(executor, + this.beanFactory, beanName)); } throw ex; } } - private Object wrapThreadPoolTaskExecutor(Object bean) { + private Object wrapThreadPoolTaskExecutor(Object bean, String beanName) { ThreadPoolTaskExecutor executor = (ThreadPoolTaskExecutor) bean; boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers()); boolean methodsFinal = anyFinalMethods(executor); boolean cglibProxy = !classFinal && !methodsFinal; - return createThreadPoolTaskExecutorProxy(bean, cglibProxy, executor); + return createThreadPoolTaskExecutorProxy(bean, cglibProxy, executor, beanName); } - private Object wrapExecutorService(Object bean) { + private Object wrapExecutorService(Object bean, String beanName) { ExecutorService executor = (ExecutorService) bean; boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers()); boolean methodFinal = anyFinalMethods(executor); boolean cglibProxy = !classFinal && !methodFinal; - return createExecutorServiceProxy(bean, cglibProxy, executor); + return createExecutorServiceProxy(bean, cglibProxy, executor, beanName); } - private Object wrapScheduledExecutorService(Object bean) { + private Object wrapScheduledExecutorService(Object bean, String beanName) { ScheduledExecutorService executor = (ScheduledExecutorService) bean; boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers()); boolean methodFinal = anyFinalMethods(executor); boolean cglibProxy = !classFinal && !methodFinal; - return createScheduledExecutorServiceProxy(bean, cglibProxy, executor); + return createScheduledExecutorServiceProxy(bean, cglibProxy, executor, beanName); } - private Object wrapAsyncTaskExecutor(Object bean) { + private Object wrapAsyncTaskExecutor(Object bean, String beanName) { AsyncTaskExecutor executor = (AsyncTaskExecutor) bean; boolean classFinal = Modifier.isFinal(bean.getClass().getModifiers()); boolean methodsFinal = anyFinalMethods(executor); boolean cglibProxy = !classFinal && !methodsFinal; - return createAsyncTaskExecutorProxy(bean, cglibProxy, executor); + return createAsyncTaskExecutorProxy(bean, cglibProxy, executor, beanName); } boolean isProxyNeeded(String beanName) { @@ -179,57 +179,62 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } Object createThreadPoolTaskExecutorProxy(Object bean, boolean cglibProxy, - ThreadPoolTaskExecutor executor) { + ThreadPoolTaskExecutor executor, String beanName) { if (!cglibProxy) { - return new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor); + return new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor, + beanName); } - return getProxiedObject(bean, cglibProxy, executor, - () -> new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor)); + return getProxiedObject(bean, beanName, true, executor, + () -> new LazyTraceThreadPoolTaskExecutor(this.beanFactory, executor, + beanName)); } Supplier createThreadPoolTaskSchedulerProxy( - ThreadPoolTaskScheduler executor) { - return () -> new LazyTraceThreadPoolTaskScheduler(this.beanFactory, executor); + ThreadPoolTaskScheduler executor, String beanName) { + return () -> new LazyTraceThreadPoolTaskScheduler(this.beanFactory, executor, + beanName); } Supplier createScheduledThreadPoolExecutorProxy( - ScheduledThreadPoolExecutor executor) { + ScheduledThreadPoolExecutor executor, String beanName) { return () -> new LazyTraceScheduledThreadPoolExecutor(executor.getCorePoolSize(), executor.getThreadFactory(), executor.getRejectedExecutionHandler(), - this.beanFactory, executor); + this.beanFactory, executor, beanName); } Object createExecutorServiceProxy(Object bean, boolean cglibProxy, - ExecutorService executor) { - return getProxiedObject(bean, cglibProxy, executor, () -> { + ExecutorService executor, String beanName) { + return getProxiedObject(bean, beanName, cglibProxy, executor, () -> { if (executor instanceof ScheduledExecutorService) { - return new TraceableScheduledExecutorService(this.beanFactory, executor); + return new TraceableScheduledExecutorService(this.beanFactory, executor, + beanName); } - - return new TraceableExecutorService(this.beanFactory, executor); + return new TraceableExecutorService(this.beanFactory, executor, beanName); }); } Object createScheduledExecutorServiceProxy(Object bean, boolean cglibProxy, - ScheduledExecutorService executor) { - return getProxiedObject(bean, cglibProxy, executor, - () -> new TraceableScheduledExecutorService(this.beanFactory, executor)); + ScheduledExecutorService executor, String beanName) { + return getProxiedObject(bean, beanName, cglibProxy, executor, + () -> new TraceableScheduledExecutorService(this.beanFactory, executor, + beanName)); } Object createAsyncTaskExecutorProxy(Object bean, boolean cglibProxy, - AsyncTaskExecutor executor) { - return getProxiedObject(bean, cglibProxy, executor, () -> { + AsyncTaskExecutor executor, String beanName) { + return getProxiedObject(bean, beanName, cglibProxy, executor, () -> { if (bean instanceof ThreadPoolTaskScheduler) { return new LazyTraceThreadPoolTaskScheduler(this.beanFactory, - (ThreadPoolTaskScheduler) executor); + (ThreadPoolTaskScheduler) executor, beanName); } - return new LazyTraceAsyncTaskExecutor(this.beanFactory, executor); + return new LazyTraceAsyncTaskExecutor(this.beanFactory, executor, beanName); }); } - private Object getProxiedObject(Object bean, boolean cglibProxy, Executor executor, - Supplier supplier) { - ProxyFactoryBean factory = proxyFactoryBean(bean, cglibProxy, executor, supplier); + private Object getProxiedObject(Object bean, String beanName, boolean cglibProxy, + Executor executor, Supplier supplier) { + ProxyFactoryBean factory = proxyFactoryBean(bean, beanName, cglibProxy, executor, + supplier); try { return getObject(factory); } @@ -246,7 +251,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { "Will wrap ThreadPoolTaskScheduler in its tracing representation due to previous errors"); } return createThreadPoolTaskSchedulerProxy( - (ThreadPoolTaskScheduler) bean).get(); + (ThreadPoolTaskScheduler) bean, beanName).get(); } else if (bean instanceof ScheduledThreadPoolExecutor) { if (log.isDebugEnabled()) { @@ -254,7 +259,7 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { "Will wrap ScheduledThreadPoolExecutor in its tracing representation due to previous errors"); } return createScheduledThreadPoolExecutorProxy( - (ScheduledThreadPoolExecutor) bean).get(); + (ScheduledThreadPoolExecutor) bean, beanName).get(); } } catch (Exception ex2) { @@ -268,17 +273,18 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { } } - private ProxyFactoryBean proxyFactoryBean(Object bean, boolean cglibProxy, - Executor executor, Supplier supplier) { + private ProxyFactoryBean proxyFactoryBean(Object bean, String beanName, + boolean cglibProxy, Executor executor, Supplier supplier) { ProxyFactoryBean factory = new ProxyFactoryBean(); factory.setProxyTargetClass(cglibProxy); - factory.addAdvice( - new ExecutorMethodInterceptor(executor, this.beanFactory) { - @Override - Executor executor(BeanFactory beanFactory, Executor executor) { - return supplier.get(); - } - }); + factory.addAdvice(new ExecutorMethodInterceptor(executor, + this.beanFactory, beanName) { + @Override + Executor executor(BeanFactory beanFactory, Executor executor, + String beanName) { + return supplier.get(); + } + }); factory.setTarget(bean); return factory; } @@ -342,14 +348,17 @@ class ExecutorMethodInterceptor implements MethodInterceptor private final BeanFactory beanFactory; - ExecutorMethodInterceptor(T delegate, BeanFactory beanFactory) { + private final String beanName; + + ExecutorMethodInterceptor(T delegate, BeanFactory beanFactory, String beanName) { this.delegate = delegate; this.beanFactory = beanFactory; + this.beanName = beanName; } @Override public Object invoke(MethodInvocation invocation) throws Throwable { - T executor = executor(this.beanFactory, this.delegate); + T executor = executor(this.beanFactory, this.delegate, this.beanName); Method methodOnTracedBean = getMethod(invocation, executor); if (methodOnTracedBean != null) { try { @@ -370,8 +379,8 @@ class ExecutorMethodInterceptor implements MethodInterceptor method.getParameterTypes()); } - T executor(BeanFactory beanFactory, T executor) { - return (T) new LazyTraceExecutor(beanFactory, executor); + T executor(BeanFactory beanFactory, T executor, String beanName) { + return (T) new LazyTraceExecutor(beanFactory, executor, beanName); } } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceAsyncTaskExecutor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceAsyncTaskExecutor.java index 3156f8fc3..9832561c8 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceAsyncTaskExecutor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceAsyncTaskExecutor.java @@ -45,6 +45,8 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor { private final AsyncTaskExecutor delegate; + private final String beanName; + private Tracing tracing; private SpanNamer spanNamer; @@ -53,13 +55,21 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor { AsyncTaskExecutor delegate) { this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = null; + } + + public LazyTraceAsyncTaskExecutor(BeanFactory beanFactory, AsyncTaskExecutor delegate, + String beanName) { + this.beanFactory = beanFactory; + this.delegate = delegate; + this.beanName = beanName; } @Override public void execute(Runnable task) { Runnable taskToRun = task; if (!ContextUtil.isContextUnusable(this.beanFactory)) { - taskToRun = new TraceRunnable(tracing(), spanNamer(), task); + taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName); } this.delegate.execute(taskToRun); } @@ -68,7 +78,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor { public void execute(Runnable task, long startTimeout) { Runnable taskToRun = task; if (!ContextUtil.isContextUnusable(this.beanFactory)) { - taskToRun = new TraceRunnable(tracing(), spanNamer(), task); + taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName); } this.delegate.execute(taskToRun, startTimeout); } @@ -77,7 +87,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor { public Future submit(Runnable task) { Runnable taskToRun = task; if (!ContextUtil.isContextUnusable(this.beanFactory)) { - taskToRun = new TraceRunnable(tracing(), spanNamer(), task); + taskToRun = new TraceRunnable(tracing(), spanNamer(), task, this.beanName); } return this.delegate.submit(taskToRun); } @@ -86,7 +96,7 @@ public class LazyTraceAsyncTaskExecutor implements AsyncTaskExecutor { public Future submit(Callable task) { Callable taskToRun = task; if (!ContextUtil.isContextUnusable(this.beanFactory)) { - taskToRun = new TraceCallable<>(tracing(), spanNamer(), task); + taskToRun = new TraceCallable<>(tracing(), spanNamer(), task, this.beanName); } return this.delegate.submit(taskToRun); } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceExecutor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceExecutor.java index 5394e66f1..6852f94a2 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceExecutor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceExecutor.java @@ -42,6 +42,8 @@ public class LazyTraceExecutor implements Executor { private final Executor delegate; + private final String beanName; + private Tracing tracing; private SpanNamer spanNamer; @@ -49,6 +51,14 @@ public class LazyTraceExecutor implements Executor { public LazyTraceExecutor(BeanFactory beanFactory, Executor delegate) { this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = null; + } + + public LazyTraceExecutor(BeanFactory beanFactory, Executor delegate, + String beanName) { + this.beanFactory = beanFactory; + this.delegate = delegate; + this.beanName = beanName; } @Override @@ -66,7 +76,8 @@ public class LazyTraceExecutor implements Executor { return; } } - this.delegate.execute(new TraceRunnable(this.tracing, spanNamer(), command)); + this.delegate.execute( + new TraceRunnable(this.tracing, spanNamer(), command, this.beanName)); } // due to some race conditions trace keys might not be ready yet diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutor.java index d2f761ab3..0c30964b0 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutor.java @@ -60,6 +60,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { private final ScheduledThreadPoolExecutor delegate; + private final String beanName; + private final Method decorateTaskRunnable; private final Method decorateTaskCallable; @@ -81,10 +83,11 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { private SpanNamer spanNamer; LazyTraceScheduledThreadPoolExecutor(int corePoolSize, BeanFactory beanFactory, - ScheduledThreadPoolExecutor delegate) { + ScheduledThreadPoolExecutor delegate, String beanName) { super(corePoolSize); this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = beanName; this.decorateTaskRunnable = ReflectionUtils.findMethod( ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class, RunnableScheduledFuture.class); @@ -122,82 +125,12 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { } LazyTraceScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory, - BeanFactory beanFactory, ScheduledThreadPoolExecutor delegate) { - super(corePoolSize, threadFactory); - this.beanFactory = beanFactory; - this.delegate = delegate; - this.decorateTaskRunnable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class, - RunnableScheduledFuture.class); - makeAccessibleIfNotNull(this.decorateTaskRunnable); - this.decorateTaskCallable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "decorateTaskCallable", Callable.class, - RunnableScheduledFuture.class); - makeAccessibleIfNotNull(this.decorateTaskCallable); - this.finalize = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "finalize"); - makeAccessibleIfNotNull(this.finalize); - this.beforeExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "beforeExecute"); - makeAccessibleIfNotNull(this.beforeExecute); - this.afterExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "afterExecute", null); - makeAccessibleIfNotNull(this.afterExecute); - this.terminated = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "terminated", null); - makeAccessibleIfNotNull(this.terminated); - this.newTaskForRunnable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "newTaskFor", Runnable.class, - Object.class); - makeAccessibleIfNotNull(this.newTaskForRunnable); - this.newTaskForCallable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "newTaskFor", Callable.class, - Object.class); - makeAccessibleIfNotNull(this.newTaskForCallable); - } - - LazyTraceScheduledThreadPoolExecutor(int corePoolSize, RejectedExecutionHandler handler, BeanFactory beanFactory, - ScheduledThreadPoolExecutor delegate) { - super(corePoolSize, handler); - this.beanFactory = beanFactory; - this.delegate = delegate; - this.decorateTaskRunnable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class, - RunnableScheduledFuture.class); - makeAccessibleIfNotNull(this.decorateTaskRunnable); - this.decorateTaskCallable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "decorateTaskCallable", Callable.class, - RunnableScheduledFuture.class); - makeAccessibleIfNotNull(this.decorateTaskCallable); - this.finalize = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "finalize", null); - makeAccessibleIfNotNull(this.finalize); - this.beforeExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "beforeExecute", null); - makeAccessibleIfNotNull(this.beforeExecute); - this.afterExecute = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "afterExecute", null); - makeAccessibleIfNotNull(this.afterExecute); - this.terminated = ReflectionUtils.findMethod(ScheduledThreadPoolExecutor.class, - "terminated", null); - makeAccessibleIfNotNull(this.terminated); - this.newTaskForRunnable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "newTaskFor", Runnable.class, - Object.class); - makeAccessibleIfNotNull(this.newTaskForRunnable); - this.newTaskForCallable = ReflectionUtils.findMethod( - ScheduledThreadPoolExecutor.class, "newTaskFor", Callable.class, - Object.class); - makeAccessibleIfNotNull(this.newTaskForCallable); - } - - LazyTraceScheduledThreadPoolExecutor(int corePoolSize, ThreadFactory threadFactory, - RejectedExecutionHandler handler, BeanFactory beanFactory, - ScheduledThreadPoolExecutor delegate) { + ScheduledThreadPoolExecutor delegate, String beanName) { super(corePoolSize, threadFactory, handler); this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = beanName; this.decorateTaskRunnable = ReflectionUtils.findMethod( ScheduledThreadPoolExecutor.class, "decorateTask", Runnable.class, RunnableScheduledFuture.class); @@ -234,7 +167,7 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { RunnableScheduledFuture task) { return (RunnableScheduledFuture) ReflectionUtils.invokeMethod( this.decorateTaskRunnable, this.delegate, - new TraceRunnable(tracing(), spanNamer(), runnable), task); + new TraceRunnable(tracing(), spanNamer(), runnable, this.beanName), task); } @Override @@ -243,57 +176,63 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { RunnableScheduledFuture task) { return (RunnableScheduledFuture) ReflectionUtils.invokeMethod( this.decorateTaskCallable, this.delegate, - new TraceCallable<>(tracing(), spanNamer(), callable), task); + new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName), + task); } @Override public ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit) { - return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), command), - delay, unit); + return this.delegate.schedule( + new TraceRunnable(tracing(), spanNamer(), command, this.beanName), delay, + unit); } @Override public ScheduledFuture schedule(Callable callable, long delay, TimeUnit unit) { return this.delegate.schedule( - new TraceCallable<>(tracing(), spanNamer(), callable), delay, unit); + new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName), + delay, unit); } @Override public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) { return this.delegate.scheduleAtFixedRate( - new TraceRunnable(tracing(), spanNamer(), command), initialDelay, period, - unit); + new TraceRunnable(tracing(), spanNamer(), command, this.beanName), + initialDelay, period, unit); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) { return this.delegate.scheduleWithFixedDelay( - new TraceRunnable(tracing(), spanNamer(), command), initialDelay, delay, - unit); + new TraceRunnable(tracing(), spanNamer(), command, this.beanName), + initialDelay, delay, unit); } @Override public void execute(Runnable command) { - this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), command)); + this.delegate.execute( + new TraceRunnable(tracing(), spanNamer(), command, this.beanName)); } @Override public Future submit(Runnable task) { - return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task)); + return this.delegate + .submit(new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public Future submit(Runnable task, T result) { - return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task), - result); + return this.delegate.submit( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), result); } @Override public Future submit(Callable task) { - return this.delegate.submit(new TraceCallable<>(tracing(), spanNamer(), task)); + return this.delegate + .submit(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName)); } @Override @@ -479,13 +418,13 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { @Override public void beforeExecute(Thread t, Runnable r) { ReflectionUtils.invokeMethod(this.beforeExecute, this.delegate, t, - new TraceRunnable(tracing(), spanNamer(), r)); + new TraceRunnable(tracing(), spanNamer(), r, this.beanName)); } @Override public void afterExecute(Runnable r, Throwable t) { ReflectionUtils.invokeMethod(this.afterExecute, this.delegate, - new TraceRunnable(tracing(), spanNamer(), r), t); + new TraceRunnable(tracing(), spanNamer(), r, this.beanName), t); } @Override @@ -497,7 +436,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { @SuppressWarnings("unchecked") public RunnableFuture newTaskFor(Runnable runnable, T value) { return (RunnableFuture) ReflectionUtils.invokeMethod(this.newTaskForRunnable, - this.delegate, new TraceRunnable(tracing(), spanNamer(), runnable), + this.delegate, + new TraceRunnable(tracing(), spanNamer(), runnable, this.beanName), value); } @@ -505,7 +445,8 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { @SuppressWarnings("unchecked") public RunnableFuture newTaskFor(Callable callable) { return (RunnableFuture) ReflectionUtils.invokeMethod(this.newTaskForCallable, - this.delegate, new TraceCallable<>(tracing(), spanNamer(), callable)); + this.delegate, + new TraceCallable<>(tracing(), spanNamer(), callable, this.beanName)); } @Override @@ -519,7 +460,7 @@ class LazyTraceScheduledThreadPoolExecutor extends ScheduledThreadPoolExecutor { List> ts = new ArrayList<>(); for (Callable task : tasks) { if (!(task instanceof TraceCallable)) { - ts.add(new TraceCallable<>(tracing(), spanNamer(), task)); + ts.add(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName)); } } return ts; diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskExecutor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskExecutor.java index 0bc3f0147..c12ab05c8 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskExecutor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskExecutor.java @@ -51,6 +51,8 @@ public class LazyTraceThreadPoolTaskExecutor extends ThreadPoolTaskExecutor { private final ThreadPoolTaskExecutor delegate; + private final String beanName; + private Tracing tracing; private SpanNamer spanNamer; @@ -59,44 +61,55 @@ public class LazyTraceThreadPoolTaskExecutor extends ThreadPoolTaskExecutor { ThreadPoolTaskExecutor delegate) { this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = null; + } + + public LazyTraceThreadPoolTaskExecutor(BeanFactory beanFactory, + ThreadPoolTaskExecutor delegate, String beanName) { + this.beanFactory = beanFactory; + this.delegate = delegate; + this.beanName = beanName; } @Override public void execute(Runnable task) { this.delegate.execute(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceRunnable(tracing(), spanNamer(), task)); + : new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public void execute(Runnable task, long startTimeout) { - this.delegate.execute(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceRunnable(tracing(), spanNamer(), task), startTimeout); + this.delegate.execute( + ContextUtil.isContextUnusable(this.beanFactory) ? task + : new TraceRunnable(tracing(), spanNamer(), task, this.beanName), + startTimeout); } @Override public Future submit(Runnable task) { return this.delegate.submit(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceRunnable(tracing(), spanNamer(), task)); + : new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public Future submit(Callable task) { return this.delegate.submit(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceCallable<>(tracing(), spanNamer(), task)); + : new TraceCallable<>(tracing(), spanNamer(), task, this.beanName)); } @Override public ListenableFuture submitListenable(Runnable task) { return this.delegate .submitListenable(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceRunnable(tracing(), spanNamer(), task)); + : new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public ListenableFuture submitListenable(Callable task) { return this.delegate .submitListenable(ContextUtil.isContextUnusable(this.beanFactory) ? task - : new TraceCallable<>(tracing(), spanNamer(), task)); + : new TraceCallable<>(tracing(), spanNamer(), task, + this.beanName)); } @Override diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskScheduler.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskScheduler.java index 113e76f5e..061f37d2b 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskScheduler.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskScheduler.java @@ -62,6 +62,8 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { private final ThreadPoolTaskScheduler delegate; + private final String beanName; + private final Method initializeExecutor; private final Method createExecutor; @@ -77,9 +79,10 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { private SpanNamer spanNamer; LazyTraceThreadPoolTaskScheduler(BeanFactory beanFactory, - ThreadPoolTaskScheduler delegate) { + ThreadPoolTaskScheduler delegate, String beanName) { this.beanFactory = beanFactory; this.delegate = delegate; + this.beanName = beanName; this.initializeExecutor = ReflectionUtils .findMethod(ThreadPoolTaskScheduler.class, "initializeExecutor", null); makeAccessibleIfNotNull(this.initializeExecutor); @@ -127,11 +130,13 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { if (executorService instanceof TraceableScheduledExecutorService) { return executorService; } - return new TraceableExecutorService(this.beanFactory, executorService); + return new TraceableExecutorService(this.beanFactory, executorService, + this.beanName); } private ThreadFactory traceThreadFactory(ThreadFactory threadFactory) { - return r -> threadFactory.newThread(new TraceRunnable(tracing(), spanNamer(), r)); + return r -> threadFactory + .newThread(new TraceRunnable(tracing(), spanNamer(), r, this.beanName)); } @Override @@ -144,14 +149,16 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { if (executorService instanceof TraceableScheduledExecutorService) { return executorService; } - return new TraceableScheduledExecutorService(this.beanFactory, executorService); + return new TraceableScheduledExecutorService(this.beanFactory, executorService, + this.beanName); } @Override public ScheduledExecutorService getScheduledExecutor() throws IllegalStateException { ScheduledExecutorService executor = this.delegate.getScheduledExecutor(); return executor instanceof TraceableScheduledExecutorService ? executor - : new TraceableScheduledExecutorService(this.beanFactory, executor); + : new TraceableScheduledExecutorService(this.beanFactory, executor, + this.beanName); } @Override @@ -164,7 +171,7 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { } return new LazyTraceScheduledThreadPoolExecutor(executor.getCorePoolSize(), executor.getThreadFactory(), executor.getRejectedExecutionHandler(), - this.beanFactory, executor); + this.beanFactory, executor, this.beanName); } @Override @@ -184,41 +191,45 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { @Override public void execute(Runnable task) { - this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), task)); + this.delegate + .execute(new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public void execute(Runnable task, long startTimeout) { - this.delegate.execute(new TraceRunnable(tracing(), spanNamer(), task), + this.delegate.execute( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTimeout); } @Override public Future submit(Runnable task) { - return this.delegate.submit(new TraceRunnable(tracing(), spanNamer(), task)); + return this.delegate + .submit(new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public Future submit(Callable task) { - return this.delegate.submit(new TraceCallable<>(tracing(), spanNamer(), task)); + return this.delegate + .submit(new TraceCallable<>(tracing(), spanNamer(), task, this.beanName)); } @Override public ListenableFuture submitListenable(Runnable task) { - return this.delegate - .submitListenable(new TraceRunnable(tracing(), spanNamer(), task)); + return this.delegate.submitListenable( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override public ListenableFuture submitListenable(Callable task) { - return this.delegate - .submitListenable(new TraceCallable<>(tracing(), spanNamer(), task)); + return this.delegate.submitListenable( + new TraceCallable<>(tracing(), spanNamer(), task, this.beanName)); } @Override public void cancelRemainingTask(Runnable task) { ReflectionUtils.invokeMethod(this.cancelRemainingTask, this.delegate, - new TraceRunnable(tracing(), spanNamer(), task)); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName)); } @Override @@ -229,13 +240,14 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { @Override @Nullable public ScheduledFuture schedule(Runnable task, Trigger trigger) { - return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task), - trigger); + return this.delegate.schedule( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), trigger); } @Override public ScheduledFuture schedule(Runnable task, Date startTime) { - return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task), + return this.delegate.schedule( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime); } @@ -243,26 +255,28 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { public ScheduledFuture scheduleAtFixedRate(Runnable task, Date startTime, long period) { return this.delegate.scheduleAtFixedRate( - new TraceRunnable(tracing(), spanNamer(), task), startTime, period); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime, + period); } @Override public ScheduledFuture scheduleAtFixedRate(Runnable task, long period) { return this.delegate.scheduleAtFixedRate( - new TraceRunnable(tracing(), spanNamer(), task), period); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), period); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable task, Date startTime, long delay) { return this.delegate.scheduleWithFixedDelay( - new TraceRunnable(tracing(), spanNamer(), task), startTime, delay); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime, + delay); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable task, long delay) { return this.delegate.scheduleWithFixedDelay( - new TraceRunnable(tracing(), spanNamer(), task), delay); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), delay); } @Override @@ -385,7 +399,8 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { @Override public ScheduledFuture schedule(Runnable task, Instant startTime) { - return this.delegate.schedule(new TraceRunnable(tracing(), spanNamer(), task), + return this.delegate.schedule( + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime); } @@ -393,26 +408,28 @@ class LazyTraceThreadPoolTaskScheduler extends ThreadPoolTaskScheduler { public ScheduledFuture scheduleAtFixedRate(Runnable task, Instant startTime, Duration period) { return this.delegate.scheduleAtFixedRate( - new TraceRunnable(tracing(), spanNamer(), task), startTime, period); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime, + period); } @Override public ScheduledFuture scheduleAtFixedRate(Runnable task, Duration period) { return this.delegate.scheduleAtFixedRate( - new TraceRunnable(tracing(), spanNamer(), task), period); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), period); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable task, Instant startTime, Duration delay) { return this.delegate.scheduleWithFixedDelay( - new TraceRunnable(tracing(), spanNamer(), task), startTime, delay); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), startTime, + delay); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable task, Duration delay) { return this.delegate.scheduleWithFixedDelay( - new TraceRunnable(tracing(), spanNamer(), task), delay); + new TraceRunnable(tracing(), spanNamer(), task, this.beanName), delay); } private Tracing tracing() { diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorService.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorService.java index abec57c2b..927e71cc2 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorService.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorService.java @@ -42,7 +42,7 @@ public class TraceableExecutorService implements ExecutorService { final ExecutorService delegate; - private final String spanName; + final String spanName; Tracing tracing; diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableScheduledExecutorService.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableScheduledExecutorService.java index 17412b60e..cbef777cd 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableScheduledExecutorService.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/TraceableScheduledExecutorService.java @@ -39,43 +39,52 @@ public class TraceableScheduledExecutorService extends TraceableExecutorService super(beanFactory, delegate); } + public TraceableScheduledExecutorService(BeanFactory beanFactory, + final ExecutorService delegate, String beanName) { + super(beanFactory, delegate, beanName); + } + private ScheduledExecutorService getScheduledExecutorService() { return (ScheduledExecutorService) this.delegate; } @Override public ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit) { - return getScheduledExecutorService().schedule( - ContextUtil.isContextUnusable(this.beanFactory) ? command - : new TraceRunnable(tracing(), spanNamer(), command), - delay, unit); + return getScheduledExecutorService() + .schedule(ContextUtil.isContextUnusable(this.beanFactory) ? command + : new TraceRunnable(tracing(), spanNamer(), command, + this.spanName), + delay, unit); } @Override public ScheduledFuture schedule(Callable callable, long delay, TimeUnit unit) { - return getScheduledExecutorService().schedule( - ContextUtil.isContextUnusable(this.beanFactory) ? callable - : new TraceCallable<>(tracing(), spanNamer(), callable), - delay, unit); + return getScheduledExecutorService() + .schedule(ContextUtil.isContextUnusable(this.beanFactory) ? callable + : new TraceCallable<>(tracing(), spanNamer(), callable, + this.spanName), + delay, unit); } @Override public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) { - return getScheduledExecutorService().scheduleAtFixedRate( - ContextUtil.isContextUnusable(this.beanFactory) ? command - : new TraceRunnable(tracing(), spanNamer(), command), - initialDelay, period, unit); + return getScheduledExecutorService() + .scheduleAtFixedRate(ContextUtil.isContextUnusable(this.beanFactory) + ? command : new TraceRunnable(tracing(), spanNamer(), command, + this.spanName), + initialDelay, period, unit); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) { - return getScheduledExecutorService().scheduleWithFixedDelay( - ContextUtil.isContextUnusable(this.beanFactory) ? command - : new TraceRunnable(tracing(), spanNamer(), command), - initialDelay, delay, unit); + return getScheduledExecutorService() + .scheduleWithFixedDelay(ContextUtil.isContextUnusable(this.beanFactory) + ? command : new TraceRunnable(tracing(), spanNamer(), command, + this.spanName), + initialDelay, delay, unit); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessorTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessorTests.java index 949dc0b07..1421dded2 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessorTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessorTests.java @@ -212,7 +212,7 @@ public class ExecutorBeanPostProcessorTests { ExecutorBeanPostProcessor bpp = new ExecutorBeanPostProcessor(this.beanFactory) { @Override Object createThreadPoolTaskExecutorProxy(Object bean, boolean cglibProxy, - ThreadPoolTaskExecutor executor) { + ThreadPoolTaskExecutor executor, String beanName) { throw new AopConfigException("foo"); } }; diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutorTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutorTests.java index 8deab8fe1..3b6008eed 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutorTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceScheduledThreadPoolExecutorTests.java @@ -39,7 +39,8 @@ public class LazyTraceScheduledThreadPoolExecutorTests { }; BeanFactory beanFactory = BDDMockito.mock(BeanFactory.class); - new LazyTraceScheduledThreadPoolExecutor(10, beanFactory, executor).finalize(); + new LazyTraceScheduledThreadPoolExecutor(10, beanFactory, executor, null) + .finalize(); BDDAssertions.then(wasCalled).isFalse(); BDDAssertions.then(executor.isShutdown()).isFalse(); diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskSchedulerTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskSchedulerTests.java index 0ed409b45..1ec8b7a26 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskSchedulerTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/LazyTraceThreadPoolTaskSchedulerTests.java @@ -58,8 +58,8 @@ public class LazyTraceThreadPoolTaskSchedulerTests { @BeforeEach public void setup() { - this.executor = new LazyTraceThreadPoolTaskScheduler(beanFactory(), - this.delegate); + this.executor = new LazyTraceThreadPoolTaskScheduler(beanFactory(), this.delegate, + null); } @AfterEach diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorServiceTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorServiceTests.java index 014b7a0f5..37c4867ec 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorServiceTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/TraceableExecutorServiceTests.java @@ -78,7 +78,7 @@ public class TraceableExecutorServiceTests { @BeforeEach public void setup() { this.traceManagerableExecutorService = new TraceableExecutorService( - beanFactory(true), this.executorService); + beanFactory(true), this.executorService, "foo"); this.spans.clear(); this.spanVerifyingRunnable.clear(); }