From 5993d0bb13371064bd7d4392453354951374a550 Mon Sep 17 00:00:00 2001 From: Francisco Javier Tirado Sarti Date: Fri, 11 Sep 2026 13:56:34 +0200 Subject: [PATCH] [Fix #1670] Stopping scheduledexecutorservice before clearing listeners Fix https://github.com/open-workflow-specification/sdk-java/issues/1670 Signed-off-by: Francisco Javier Tirado Sarti --- .../impl/DefaultExecutorServiceFactory.java | 6 +----- .../impl/WorkflowApplication.java | 6 ++---- .../io/serverlessworkflow/impl/WorkflowUtils.java | 14 ++++++++++++++ 3 files changed, 17 insertions(+), 9 deletions(-) diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java b/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java index 8d43b7d7d..8b882c143 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java @@ -17,7 +17,6 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; public class DefaultExecutorServiceFactory implements ExecutorServiceFactory { private ExecutorService service = Executors.newCachedThreadPool(); @@ -29,9 +28,6 @@ public ExecutorService get() { @Override public void close() throws Exception { - if (!service.isShutdown()) { - service.shutdown(); - service.awaitTermination(2, TimeUnit.SECONDS); - } + WorkflowUtils.safeShutdown(service); } } diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java index 59b3dbe4b..73367d438 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java @@ -17,6 +17,7 @@ import static io.serverlessworkflow.impl.WorkflowUtils.loadFirst; import static io.serverlessworkflow.impl.WorkflowUtils.safeClose; +import static io.serverlessworkflow.impl.WorkflowUtils.safeShutdown; import io.serverlessworkflow.api.types.SchemaInline; import io.serverlessworkflow.api.types.Workflow; @@ -585,11 +586,11 @@ public WorkflowDefinition workflowDefinition(Workflow workflow) { @Override public void close() { safeClose(executorFactory); + safeShutdown(schedulerExecutorService); for (EventPublisher eventPublisher : eventPublishers) { safeClose(eventPublisher); } safeClose(eventConsumer); - for (WorkflowDefinition definition : definitions.values()) { safeClose(definition); } @@ -604,9 +605,6 @@ public void close() { } listenersByPriority.clear(); } - if (this.schedulerExecutorService != null) { - schedulerExecutorService.shutdownNow(); - } } public WorkflowPositionFactory positionFactory() { diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java index f1e3e35f4..b035c52c2 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java @@ -40,6 +40,8 @@ import java.util.Objects; import java.util.Optional; import java.util.ServiceLoader; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.TimeUnit; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -177,6 +179,18 @@ public static void safeClose(AutoCloseable closeable) { } } + public static void safeShutdown(ExecutorService service) { + if (service != null && !service.isShutdown()) { + try { + service.shutdownNow(); + service.awaitTermination(2, TimeUnit.SECONDS); + } catch (InterruptedException ex) { + logger.warn("Thread was interrupted when awaiting service task termination", ex); + Thread.currentThread().interrupt(); + } + } + } + public static boolean whenExceptTest( Optional whenFilter, Optional exceptFilter,