From e0e572a8b4637f5ba9d830425fbf9849bec5d9fd Mon Sep 17 00:00:00 2001 From: Noemi Szemenyei Date: Fri, 11 Mar 2022 09:01:25 +0100 Subject: [PATCH] update (cherry picked from commit dab8ed8a3b44966bcc9f250ba0c70e822c60ace7) --- .../config/BackgroundExecutorWrapper.java | 19 ++++ .../executor/DefaultBackgroundExecutor.java | 88 +++++++++---------- .../org/tests/cache/TestBeanCacheAsync.java | 51 +++++------ 3 files changed, 85 insertions(+), 73 deletions(-) diff --git a/ebean-api/src/main/java/io/ebean/config/BackgroundExecutorWrapper.java b/ebean-api/src/main/java/io/ebean/config/BackgroundExecutorWrapper.java index 289b7327a..306b7e88e 100644 --- a/ebean-api/src/main/java/io/ebean/config/BackgroundExecutorWrapper.java +++ b/ebean-api/src/main/java/io/ebean/config/BackgroundExecutorWrapper.java @@ -20,4 +20,23 @@ public interface BackgroundExecutorWrapper { */ Runnable wrap(Runnable task); + /** + * Combines two wrappers by joining them. + */ + default BackgroundExecutorWrapper with(BackgroundExecutorWrapper inner) { + return new BackgroundExecutorWrapper() { + + @Override + public Runnable wrap(Runnable task) { + return BackgroundExecutorWrapper.this.wrap(inner.wrap(task)); + } + + @Override + public Callable wrap(Callable task) { + return BackgroundExecutorWrapper.this.wrap(inner.wrap(task)); + } + }; + + } + } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/executor/DefaultBackgroundExecutor.java b/ebean-core/src/main/java/io/ebeaninternal/server/executor/DefaultBackgroundExecutor.java index ca1f18947..468f397ff 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/executor/DefaultBackgroundExecutor.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/executor/DefaultBackgroundExecutor.java @@ -36,9 +36,9 @@ public final class DefaultBackgroundExecutor implements SpiBackgroundExecutor { */ Callable wrap(Callable task) { if (wrapper == null) { - return clock(task); + return task; } else { - return wrapper.wrap(clock(task)); + return wrapper.wrap(task); } } @@ -47,41 +47,43 @@ public final class DefaultBackgroundExecutor implements SpiBackgroundExecutor { */ Runnable wrap(Runnable task) { if (wrapper == null) { - return clock(task); + return task; } else { - return wrapper.wrap(clock(task)); + return wrapper.wrap(task); } } - private Callable clock(Callable task) { - if (logger.isTraceEnabled()) { - long queued = System.nanoTime(); - logger.trace("Queued {}", task); - return () -> { - long start = System.nanoTime(); - logger.trace("Start {} (delay time {} us)", task, (start - queued) / 1000L); - T ret = task.call(); - logger.trace("Stop {} (exec time {} us)", task, (System.nanoTime() - start) / 1000L); - return ret; - }; - } else { - return task; - } - } - - private Runnable clock(Runnable task) { - if (logger.isTraceEnabled()) { - long queued = System.nanoTime(); - logger.trace("Queued {}", task); - return () -> { - long start = System.nanoTime(); - logger.trace("Start {} (delay time {} us)", task, (start - queued) / 1000L); - task.run(); - logger.trace("Stop {} (exec time {} us)", task, (System.nanoTime() - start) / 1000L); - }; - } else { - return task; - } + + /** + * Decorates a runnable by adding an exception handler and some timing metrics. + * This is used in methods that accepts a Runnable and return + * either void or ScheduledFuture, as there is + * normally no Future.get() call. + * + * Note: When submitting a Callable, you must check + * Future.get() for exceptions. + */ + private Runnable logExceptions(Runnable task) { + long queued = System.nanoTime(); + logger.trace("Queued {}", task); + return () -> { + try { + if (logger.isTraceEnabled()) { + long start = System.nanoTime(); + logger.trace("Start {} (delay time {} us)", task, (start - queued) / 1000L); + task.run(); + logger.trace("Stop {} (exec time {} us)", task, (System.nanoTime() - start) / 1000L); + } else { + task.run(); + } + } catch (Throwable t) { + // log any exception here. Note they will not bubble up to the calling user + // unless Future.get() is checked. (Which is almost never done on scheduled + // background executions) + logger.error("Error while executing the task {}", task, t); + throw t; + } + }; } @Override @@ -97,44 +99,40 @@ public final class DefaultBackgroundExecutor implements SpiBackgroundExecutor { return pool.submit(wrap(task)); } + @Override public void execute(Runnable task) { - submit(() -> { - try { - task.run(); - } catch (Throwable t) { - logger.error("Error while executing the task {}", task, t); - } - }); + submit(logExceptions(task)); } @Override public void executePeriodically(Runnable task, long delay, TimeUnit unit) { - schedulePool.scheduleWithFixedDelay(wrap(task), delay, delay, unit); + schedulePool.scheduleWithFixedDelay(wrap(logExceptions(task)), delay, delay, unit); } @Override public void executePeriodically(Runnable task, long initialDelay, long delay, TimeUnit unit) { - schedulePool.scheduleWithFixedDelay(wrap(task), initialDelay, delay, unit); + schedulePool.scheduleWithFixedDelay(wrap(logExceptions(task)), initialDelay, delay, unit); } @Override public ScheduledFuture scheduleWithFixedDelay(Runnable task, long initialDelay, long delay, TimeUnit unit) { - return schedulePool.scheduleWithFixedDelay(wrap(task), initialDelay, delay, unit); + return schedulePool.scheduleWithFixedDelay(wrap(logExceptions(task)), initialDelay, delay, unit); } @Override public ScheduledFuture scheduleAtFixedRate(Runnable task, long initialDelay, long delay, TimeUnit unit) { - return schedulePool.scheduleAtFixedRate(wrap(task), initialDelay, delay, unit); + return schedulePool.scheduleAtFixedRate(wrap(logExceptions(task)), initialDelay, delay, unit); } @Override public ScheduledFuture schedule(Runnable task, long delay, TimeUnit unit) { - return schedulePool.schedule(wrap(task), delay, unit); + return schedulePool.schedule(wrap(logExceptions(task)), delay, unit); } @Override public ScheduledFuture schedule(Callable task, long delay, TimeUnit unit) { + // Note: Here is no "logExceptions", because it is intended to check Future.get() by the invoker return schedulePool.schedule(wrap(task), delay, unit); } diff --git a/ebean-test/src/test/java/org/tests/cache/TestBeanCacheAsync.java b/ebean-test/src/test/java/org/tests/cache/TestBeanCacheAsync.java index 41430f246..7dd1303a4 100644 --- a/ebean-test/src/test/java/org/tests/cache/TestBeanCacheAsync.java +++ b/ebean-test/src/test/java/org/tests/cache/TestBeanCacheAsync.java @@ -12,6 +12,7 @@ import io.ebean.BaseTestCase; import io.ebean.DB; import io.ebean.Database; import io.ebean.DatabaseFactory; +import io.ebean.config.BackgroundExecutorWrapper; import io.ebean.config.CurrentTenantProvider; import io.ebean.config.DatabaseConfig; import io.ebean.config.MdcBackgroundExecutorWrapper; @@ -31,42 +32,35 @@ public class TestBeanCacheAsync extends BaseTestCase { return tenantId.get(); } } + /** * Copy tenant info to the background thread. */ - class TenantCopyBackgroundExecutorWrapper extends MdcBackgroundExecutorWrapper { + class TenantCopyBackgroundExecutorWrapper implements BackgroundExecutorWrapper { @Override public Callable wrap(Callable task) { - String tenant = tenantId.get(); - if (tenant == null) { - return super.wrap(task); - } else { - return () -> { - tenantId.set(tenant); - try { - return super.wrap(task).call(); - } finally { - tenantId.remove(); - } - }; - } + String tenant = tenantId.get(); // executed in current thread + return () -> { + tenantId.set(tenant); // executed in other thread + try { + return task.call(); + } finally { + tenantId.remove(); + } + }; } - + @Override public Runnable wrap(Runnable task) { String tenant = tenantId.get(); - if (tenant == null) { - return super.wrap(task); - } else { - return () -> { - tenantId.set(tenant); - try { - super.wrap(task).run(); - } finally { - tenantId.remove(); - } - }; - } + return () -> { + tenantId.set(tenant); + try { + task.run(); + } finally { + tenantId.remove(); + } + }; } } @@ -84,7 +78,8 @@ public class TestBeanCacheAsync extends BaseTestCase { config.setRegister(false); config.setServerCachePlugin(new DefaultServerCachePlugin()); // disables foreground local caching (as it is done in Hz/Ignite) config.setCurrentTenantProvider(new ThreadLocalTenantProvider()); - config.setBackgroundExecutorWrapper(new TenantCopyBackgroundExecutorWrapper()); + config.setBackgroundExecutorWrapper( + new MdcBackgroundExecutorWrapper().with(new TenantCopyBackgroundExecutorWrapper())); tenantId.set("4711"); Database db = DatabaseFactory.create(config);