Skip to content

Commit 802ca19

Browse files
committed
[Fix #1657] Integrate status change events with others
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 eb120fe commit 802ca19

8 files changed

Lines changed: 108 additions & 108 deletions

File tree

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

Lines changed: 1 addition & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -17,23 +17,12 @@
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;
2220

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

2724
@Override
2825
public ExecutorService get() {
29-
try {
30-
serviceLock.lock();
31-
if (service == null) {
32-
service = Executors.newCachedThreadPool();
33-
}
34-
} finally {
35-
serviceLock.unlock();
36-
}
3726
return service;
3827
}
3928
}

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)

impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java

Lines changed: 67 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ protected WorkflowMutableInstance(WorkflowDefinition definition, String id, Work
7373
public CompletableFuture<WorkflowModel> start() {
7474
return startExecution(
7575
() -> {
76+
status(WorkflowStatus.RUNNING);
7677
startedAt = Instant.now();
7778
return publishEvent(
7879
workflowContext, l -> l.onWorkflowStarted(new WorkflowStartedEvent(workflowContext)));
@@ -85,8 +86,6 @@ protected final CompletableFuture<WorkflowModel> startExecution(
8586
if (future != null) {
8687
return future;
8788
}
88-
status(WorkflowStatus.RUNNING);
89-
9089
future =
9190
runnable
9291
.get()
@@ -103,20 +102,23 @@ protected final CompletableFuture<WorkflowModel> startExecution(
103102
.orElse(input))
104103
.whenComplete(this::setCompleteDate)
105104
.thenApply(this::filterAndValidate)
106-
.whenComplete(this::handleException)
107-
.thenCompose(
108-
model ->
109-
publishEvent(
110-
workflowContext,
111-
l ->
112-
l.onWorkflowCompleted(
113-
new WorkflowCompletedEvent(workflowContext, model)))
114-
.thenApply(__ -> model)))
105+
.thenCompose(this::publishEvents)
106+
.exceptionallyCompose(this::handleException))
115107
.whenComplete(this::cleanUp);
116108
futureRef.set(future);
117109
return future;
118110
}
119111

112+
private CompletableFuture<WorkflowModel> publishEvents(WorkflowModel model) {
113+
return status(WorkflowStatus.COMPLETED)
114+
.thenCompose(
115+
__ ->
116+
publishEvent(
117+
workflowContext,
118+
l -> l.onWorkflowCompleted(new WorkflowCompletedEvent(workflowContext, model))))
119+
.thenApply(__ -> model);
120+
}
121+
120122
private void setCompleteDate(WorkflowModel result, Throwable ex) {
121123
completedAt = Instant.now();
122124
}
@@ -130,19 +132,19 @@ private void cleanUp(WorkflowModel result, Throwable ex) {
130132
workflowContext.definition().removeInstance(this);
131133
}
132134

133-
private void handleException(WorkflowModel result, Throwable exception) {
134-
if (exception != null) {
135-
final Throwable cause =
136-
exception instanceof CompletionException ? exception.getCause() : exception;
137-
if (!(cause instanceof CancellationException)) {
138-
status(WorkflowStatus.FAULTED);
139-
publishEvent(
140-
workflowContext,
141-
l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, cause)));
142-
}
143-
} else {
144-
status(WorkflowStatus.COMPLETED);
135+
private CompletableFuture<WorkflowModel> handleException(Throwable exception) {
136+
final Throwable cause =
137+
exception instanceof CompletionException ? exception.getCause() : exception;
138+
if (!(cause instanceof CancellationException)) {
139+
return status(WorkflowStatus.FAULTED)
140+
.thenCompose(
141+
__ ->
142+
publishEvent(
143+
workflowContext,
144+
l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, cause))))
145+
.thenCompose(__ -> CompletableFuture.failedFuture(exception));
145146
}
147+
return CompletableFuture.failedFuture(exception);
146148
}
147149

148150
private WorkflowModel filterAndValidate(WorkflowModel model) {
@@ -208,15 +210,19 @@ public <T> T outputAs(Class<T> clazz) {
208210
: null;
209211
}
210212

211-
public void status(WorkflowStatus state) {
213+
public CompletableFuture<?> status(WorkflowStatus state) {
212214
WorkflowStatus prevState = this.status.getAndSet(state);
213-
if (prevState != state) {
214-
publishEvent(
215-
workflowContext,
216-
l ->
217-
l.onWorkflowStatusChanged(
218-
new WorkflowStatusEvent(workflowContext, prevState, state)));
219-
}
215+
return prevState != state
216+
? publishEvent(
217+
workflowContext,
218+
l ->
219+
l.onWorkflowStatusChanged(
220+
new WorkflowStatusEvent(workflowContext, prevState, state)))
221+
: CompletableFuture.completedFuture(null);
222+
}
223+
224+
protected final void setStatus(WorkflowStatus state) {
225+
this.status.set(state);
220226
}
221227

222228
@Override
@@ -234,8 +240,9 @@ public String toString() {
234240

235241
@Override
236242
public boolean suspend() {
237-
boolean result = _suspend();
243+
boolean result = internalSuspend();
238244
if (result) {
245+
status(WorkflowStatus.SUSPENDED);
239246
publishEvent(
240247
workflowContext, l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)));
241248
}
@@ -244,19 +251,22 @@ public boolean suspend() {
244251

245252
@Override
246253
public CompletableFuture<Boolean> suspendFuture() {
247-
return _suspend()
248-
? publishEvent(
249-
workflowContext,
250-
l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)))
254+
return internalSuspend()
255+
? status(WorkflowStatus.SUSPENDED)
256+
.thenCompose(
257+
__ ->
258+
publishEvent(
259+
workflowContext,
260+
l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))))
251261
.thenApply(__ -> true)
252262
: CompletableFuture.completedFuture(false);
253263
}
254264

255-
private boolean _suspend() {
265+
private boolean internalSuspend() {
256266
try {
257267
statusLock.lock();
258268
if (TaskExecutorHelper.isActive(status.get()) && suspended == null) {
259-
internalSuspend();
269+
setSuspended();
260270
return true;
261271
} else {
262272
return false;
@@ -266,14 +276,13 @@ private boolean _suspend() {
266276
}
267277
}
268278

269-
protected final void internalSuspend() {
279+
protected final void setSuspended() {
270280
suspended = new ConcurrentHashMap<>();
271-
status(WorkflowStatus.SUSPENDED);
272281
}
273282

274283
@Override
275284
public boolean resume() {
276-
boolean result = _resume();
285+
boolean result = internalResume();
277286
if (result) {
278287
publishEvent(
279288
workflowContext, l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)));
@@ -283,15 +292,15 @@ public boolean resume() {
283292

284293
@Override
285294
public CompletableFuture<Boolean> resumeFuture() {
286-
return _resume()
295+
return internalResume()
287296
? publishEvent(
288297
workflowContext,
289298
l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)))
290299
.thenApply(__ -> true)
291300
: CompletableFuture.completedFuture(false);
292301
}
293302

294-
private boolean _resume() {
303+
private boolean internalResume() {
295304
boolean result;
296305
try {
297306
statusLock.lock();
@@ -327,25 +336,28 @@ public CompletableFuture<TaskContext> cancelCheck(TaskContext t) {
327336
}
328337

329338
public CompletableFuture<TaskContext> suspendedCheck(TaskContext t) {
339+
boolean isActive;
330340
try {
331341
statusLock.lock();
332342
if (suspended != null) {
333343
CompletableFuture<TaskContext> suspendedTask = new CompletableFuture<TaskContext>();
334344
suspended.put(suspendedTask, t);
335345
return suspendedTask;
336-
} else if (TaskExecutorHelper.isActive(status.get())) {
337-
status(WorkflowStatus.RUNNING);
338346
}
347+
isActive = TaskExecutorHelper.isActive(status.get());
339348
} finally {
340349
statusLock.unlock();
341350
}
342-
return CompletableFuture.completedFuture(t);
351+
return isActive
352+
? status(WorkflowStatus.RUNNING).thenApply(__ -> t)
353+
: CompletableFuture.completedFuture(t);
343354
}
344355

345356
@Override
346357
public boolean cancel() {
347-
boolean result = _cancel();
358+
boolean result = internalCancel();
348359
if (result) {
360+
status(WorkflowStatus.CANCELLED);
349361
publishEvent(
350362
workflowContext, l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext)));
351363
}
@@ -354,23 +366,25 @@ public boolean cancel() {
354366

355367
@Override
356368
public CompletableFuture<Boolean> cancelFuture() {
357-
return _cancel()
358-
? publishEvent(
359-
workflowContext,
360-
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext)))
369+
return internalCancel()
370+
? status(WorkflowStatus.CANCELLED)
371+
.thenCompose(
372+
__ ->
373+
publishEvent(
374+
workflowContext,
375+
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))))
361376
.thenApply(__ -> true)
362377
: CompletableFuture.completedFuture(false);
363378
}
364379

365-
private boolean _cancel() {
380+
private boolean internalCancel() {
366381
boolean result;
367382
Collection<CompletableFuture<?>> toCancel = null;
368383
try {
369384
statusLock.lock();
370385
if (TaskExecutorHelper.isActive(status.get())) {
371386
toCancel = new ArrayList<>(cancelables);
372387
cancelables.clear();
373-
status(WorkflowStatus.CANCELLED);
374388
result = true;
375389
} else {
376390
result = false;

impl/core/src/main/java/io/serverlessworkflow/impl/executors/AbstractTaskExecutor.java

Lines changed: 19 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,7 @@ private CompletableFuture<TaskContext> executeNext(
193193
WorkflowContext workflow, TaskContext taskContext) {
194194
TransitionInfo transition = taskContext.transition();
195195
if (transition.isEndNode()) {
196-
workflow.instance().status(WorkflowStatus.COMPLETED);
196+
return workflow.instance().status(WorkflowStatus.COMPLETED).thenApply(__ -> taskContext);
197197
} else if (transition.next() != null) {
198198
return transition.next().apply(workflow, taskContext.parent(), taskContext.output());
199199
}
@@ -241,15 +241,6 @@ public CompletableFuture<TaskContext> apply(
241241
})
242242
.thenCompose(t -> execute(workflowContext, t))
243243
.thenCompose(workflowContext.instance()::cancelCheck)
244-
.whenComplete(
245-
(t, e) -> {
246-
if (e != null) {
247-
handleException(
248-
workflowContext,
249-
taskContext,
250-
e instanceof CompletionException ? e.getCause() : e);
251-
}
252-
})
253244
.thenApply(
254245
t -> {
255246
outputProcessor.ifPresent(
@@ -272,7 +263,13 @@ public CompletableFuture<TaskContext> apply(
272263
l ->
273264
l.onTaskCompleted(
274265
new TaskCompletedEvent(workflowContext, taskContext)))
275-
.thenApply(__ -> t));
266+
.thenApply(__ -> t))
267+
.exceptionallyCompose(
268+
e ->
269+
handleException(
270+
workflowContext,
271+
taskContext,
272+
e instanceof CompletionException ? e.getCause() : e));
276273
if (timeout.isPresent()) {
277274
completable =
278275
completable
@@ -298,17 +295,21 @@ public CompletableFuture<TaskContext> apply(
298295
}
299296
}
300297

301-
private void handleException(
298+
private CompletableFuture<TaskContext> handleException(
302299
WorkflowContext workflowContext, TaskContext taskContext, Throwable e) {
300+
CompletableFuture<?> events;
303301
if (e instanceof CancellationException) {
304-
publishEvent(
305-
workflowContext,
306-
l -> l.onTaskCancelled(new TaskCancelledEvent(workflowContext, taskContext)));
302+
events =
303+
publishEvent(
304+
workflowContext,
305+
l -> l.onTaskCancelled(new TaskCancelledEvent(workflowContext, taskContext)));
307306
} else {
308-
publishEvent(
309-
workflowContext,
310-
l -> l.onTaskFailed(new TaskFailedEvent(workflowContext, taskContext, e)));
307+
events =
308+
publishEvent(
309+
workflowContext,
310+
l -> l.onTaskFailed(new TaskFailedEvent(workflowContext, taskContext, e)));
311311
}
312+
return events.thenCompose(__ -> CompletableFuture.failedFuture(e));
312313
}
313314

314315
public WorkflowPosition position() {

0 commit comments

Comments
 (0)