From 3e5b41928b3a940681a7090afe64619b513876b2 Mon Sep 17 00:00:00 2001 From: fjtirado Date: Fri, 10 Jul 2026 16:35:32 +0200 Subject: [PATCH] [Fix #1529] Printing error message Also improve the time spent by LifeCycleEventTest Signed-off-by: fjtirado --- .../scheduler/ScheduledInstanceRunnable.java | 22 ++++++- .../impl/test/LifeCycleEventsTest.java | 61 ++++++++----------- .../resources/workflows-samples/wait-set.yaml | 2 +- 3 files changed, 49 insertions(+), 36 deletions(-) diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/scheduler/ScheduledInstanceRunnable.java b/impl/core/src/main/java/io/serverlessworkflow/impl/scheduler/ScheduledInstanceRunnable.java index ef75af05d..d77f9b302 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/scheduler/ScheduledInstanceRunnable.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/scheduler/ScheduledInstanceRunnable.java @@ -19,9 +19,13 @@ import io.serverlessworkflow.impl.WorkflowInstance; import io.serverlessworkflow.impl.WorkflowModel; import java.util.function.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; public class ScheduledInstanceRunnable implements Runnable, Consumer { + private static final Logger logger = LoggerFactory.getLogger(ScheduledInstanceRunnable.class); + protected final WorkflowDefinition definition; public ScheduledInstanceRunnable(WorkflowDefinition definition) { @@ -45,6 +49,22 @@ public static void runScheduledInstance(WorkflowDefinition definition, WorkflowM private static void runScheduledInstance( WorkflowDefinition definition, WorkflowInstance instance) { definition.addScheduledInstance(instance); - definition.application().executorService().execute(() -> instance.start()); + definition + .application() + .executorService() + .execute( + () -> + instance + .start() + .whenComplete( + (model, ex) -> { + if (ex != null) { + logger.error( + "Scheduled workflow instance {} of definition {} has failed", + instance.id(), + definition.id(), + ex); + } + })); } } diff --git a/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java b/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java index 9a841abc7..892be2950 100644 --- a/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java +++ b/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java @@ -25,6 +25,7 @@ import io.serverlessworkflow.api.types.Workflow; import io.serverlessworkflow.fluent.spec.WorkflowBuilder; import io.serverlessworkflow.fluent.spec.dsl.DSL; +import io.serverlessworkflow.impl.ExecutorServiceFactory; import io.serverlessworkflow.impl.WorkflowApplication; import io.serverlessworkflow.impl.WorkflowDefinition; import io.serverlessworkflow.impl.WorkflowDefinitionId; @@ -54,6 +55,8 @@ import java.util.concurrent.CompletionException; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.stream.Stream; @@ -76,9 +79,25 @@ static void init() { for (String type : AbstractLifeCyclePublisher.getLifeCycleTypes()) { eventBroker.register(type, ce -> publishedEvents.add(ce)); } + ; appl = WorkflowApplication.builder() .withLifeCycleCloudEventFactory(new InputOutputLifeCycleCloudEventFactory()) + .withExecutorFactory( + new ExecutorServiceFactory() { + + private ExecutorService service = Executors.newFixedThreadPool(2); + + @Override + public void close() throws Exception { + service.shutdownNow(); + } + + @Override + public ExecutorService get() { + return service; + } + }) .withEventConsumer(eventBroker) .withEventPublisher(eventBroker) .build(); @@ -158,34 +177,6 @@ private void doTestSuspendResumeNotWait(Workflow workflow) .isBeforeOrEqualTo(workflowResumedEvent.resumedAt()); } - @ParameterizedTest(name = "{0}") - @MethodSource("waitSetWorkflowSources") - void testSuspendResumeWait(String sourceName, Workflow workflow) - throws ExecutionException, InterruptedException, TimeoutException { - doTestSuspendResumeWait(workflow); - } - - private void doTestSuspendResumeWait(Workflow workflow) - throws InterruptedException, ExecutionException, TimeoutException { - WorkflowDefinition def = appl.workflowDefinition(workflow); - WorkflowInstance instance = def.instance(Map.of()); - CompletableFuture future = instance.start(); - assertThat(instance.status()).isEqualTo(WorkflowStatus.WAITING); - instance.suspend(); - assertThat(instance.status()).isEqualTo(WorkflowStatus.SUSPENDED); - Thread.sleep(550); - instance.resume(); - assertThat(future.get(1, TimeUnit.SECONDS).asMap().orElseThrow()) - .isEqualTo(Map.of("name", "Javierito")); - assertThat(instance.status()).isEqualTo(WorkflowStatus.COMPLETED); - WorkflowSuspendedCEData workflowSuspendedEvent = - assertPojoInCE( - "io.serverlessworkflow.workflow.suspended.v1", WorkflowSuspendedCEData.class); - WorkflowResumedCEData workflowResumedEvent = - assertPojoInCE("io.serverlessworkflow.workflow.resumed.v1", WorkflowResumedCEData.class); - assertThat(workflowSuspendedEvent.suspendedAt()).isBefore(workflowResumedEvent.resumedAt()); - } - @ParameterizedTest(name = "{0}") @MethodSource("waitSetWorkflowSources") void testCancel(String sourceName, Workflow workflow) throws IOException { @@ -206,16 +197,19 @@ private void doTestCancel(Workflow workflow) { @ParameterizedTest(name = "{0}") @MethodSource("waitSetWorkflowSources") - void testSuspendResumeTimeout(String sourceName, Workflow workflow) { - doTestSuspendResumeTimeout(workflow); + void testSuspendTimeout(String sourceName, Workflow workflow) { + doTestSuspendTimeout(workflow); } - private static void doTestSuspendResumeTimeout(Workflow workflow) { + private static void doTestSuspendTimeout(Workflow workflow) { WorkflowDefinition def = appl.workflowDefinition(workflow); WorkflowInstance instance = def.instance(Map.of()); CompletableFuture future = instance.start(); instance.suspend(); - assertThat(catchThrowableOfType(TimeoutException.class, () -> future.get(1, TimeUnit.SECONDS))) + assertThat(instance.status()).isEqualTo(WorkflowStatus.SUSPENDED); + assertThat( + catchThrowableOfType( + TimeoutException.class, () -> future.get(400, TimeUnit.MILLISECONDS))) .isNotNull(); } @@ -255,11 +249,10 @@ private T assertPojoInCE(String type, Class clazz) { private static Workflow waitTestWorkflow() { return WorkflowBuilder.workflow("wait-test-java-dsl", "test", "0.1.0") .tasks( - // wait 500 ms DSL.wait( "waitABit", timeoutBuilder -> - timeoutBuilder.duration(durationBuilder -> durationBuilder.milliseconds(500))), + timeoutBuilder.duration(durationBuilder -> durationBuilder.milliseconds(200))), DSL.set("useExpression", setTaskBuilder -> setTaskBuilder.put("name", "Javierito"))) .build(); } diff --git a/impl/test/src/test/resources/workflows-samples/wait-set.yaml b/impl/test/src/test/resources/workflows-samples/wait-set.yaml index 0fe9c4fab..f3d43a57a 100644 --- a/impl/test/src/test/resources/workflows-samples/wait-set.yaml +++ b/impl/test/src/test/resources/workflows-samples/wait-set.yaml @@ -6,7 +6,7 @@ document: do: - waitABit: wait: - milliseconds: 500 + milliseconds: 200 - useExpression: set: name : Javierito \ No newline at end of file