Skip to content

Commit c481ddb

Browse files
authored
[Fix #1657] Integrate status change events with others (#1663)
The idea is to not continue workflow execution till the listeners are completed. Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com>
1 parent 17cda82 commit c481ddb

10 files changed

Lines changed: 210 additions & 222 deletions

File tree

‎impl/core/src/main/java/io/serverlessworkflow/impl/AbstractExecutorServiceHolder.java‎

Lines changed: 0 additions & 32 deletions
This file was deleted.

‎impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java‎

Lines changed: 11 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -17,23 +17,21 @@
1717

1818
import java.util.concurrent.ExecutorService;
1919
import java.util.concurrent.Executors;
20-
import java.util.concurrent.locks.Lock;
21-
import java.util.concurrent.locks.ReentrantLock;
20+
import java.util.concurrent.TimeUnit;
2221

23-
public class DefaultExecutorServiceFactory extends AbstractExecutorServiceHolder {
24-
25-
private Lock serviceLock = new ReentrantLock();
22+
public class DefaultExecutorServiceFactory implements ExecutorServiceFactory {
23+
private ExecutorService service = Executors.newCachedThreadPool();
2624

2725
@Override
2826
public ExecutorService get() {
29-
try {
30-
serviceLock.lock();
31-
if (service == null) {
32-
service = Executors.newCachedThreadPool();
33-
}
34-
} finally {
35-
serviceLock.unlock();
36-
}
3727
return service;
3828
}
29+
30+
@Override
31+
public void close() throws Exception {
32+
if (!service.isShutdown()) {
33+
service.shutdown();
34+
service.awaitTermination(2, TimeUnit.SECONDS);
35+
}
36+
}
3937
}

‎impl/core/src/main/java/io/serverlessworkflow/impl/ExecutorServiceHolder.java‎

Lines changed: 0 additions & 30 deletions
This file was deleted.

‎impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -254,7 +254,7 @@ public SchemaValidator getValidator(SchemaInline inline) {
254254
private WorkflowPositionFactory positionFactory = () -> new QueueWorkflowPosition();
255255
private WorkflowInstanceIdFactory idFactory;
256256
private WorkflowScheduler scheduler;
257-
private ExecutorServiceFactory executorFactory = new DefaultExecutorServiceFactory();
257+
private ExecutorServiceFactory executorFactory;
258258
private EventConsumer<?, ?> eventConsumer;
259259
private Collection<EventPublisher> eventPublishers = new ArrayList<>();
260260
private RuntimeDescriptorFactory descriptorFactory =
@@ -458,7 +458,9 @@ public Builder withCronResolverFactory(CronResolverFactory cronResolverFactory)
458458
}
459459

460460
public WorkflowApplication build() {
461-
461+
if (executorFactory == null) {
462+
executorFactory = new DefaultExecutorServiceFactory();
463+
}
462464
if (modelFactory == null) {
463465
modelFactory =
464466
loadFirst(WorkflowModelFactory.class)

0 commit comments

Comments
 (0)