Skip to content

Commit c8f334a

Browse files
edeandreafjtirado
andauthored
[Fix #1749] Publish onWorkflowCancelled before cancelling pending futures (#1751)
* [Fix #1749] Publish onWorkflowCancelled before cancelling pending futures Cancelling the futures registered through addCancelable (e.g. by a listen task) completes the execution pipeline synchronously, which runs cleanUp and clears the instance metadata. cancel()/cancelFuture() used to do that before publishing the status change and onWorkflowCancelled, so listeners could not find their per-instance metadata when the cancelled event arrived. internalCancel() now returns the futures to cancel, and they are cancelled only after the cancelled status change and onWorkflowCancelled have been published. Signed-off-by: Eric Deandrea <eric@ericdeandrea.dev> * [Fix #1749] Keep cancellation final and clean up after it is published Address review feedback on the cancellation ordering: - status(WorkflowStatus) now checks and sets the status under statusLock and ignores changes once the instance is CANCELLED, so a listen task receiving an event while the cancellation is being published can no longer move the instance back to WAITING and let it complete normally. - The execution pipeline waits for the cancellation to be published before running cleanUp, so instance metadata is still available to onWorkflowCancelled listeners even when the pipeline ends on its own (an event completes the listen task, or a wait elapses) while an asynchronous status change listener is still pending. Adds regressions that hold the CANCELLED status change pending and deliver a late event or let the pipeline finish in the meantime. Signed-off-by: Eric Deandrea <eric@ericdeandrea.dev> * [Fix #1749] Fail terminal transitions of a cancelled instance status(WorkflowStatus) ignored every change once the instance was CANCELLED by returning a normally completed future, but publishEvents() and handleException() ignore that value. If the cancellation happened while the last task's onTaskCompleted/onTaskFailed listener was still pending, releasing it published onWorkflowCompleted (and completed start() normally) or onWorkflowFailed for a cancelled instance. A rejected COMPLETED or FAULTED transition now fails with a CancellationException, so the pipeline ends as cancelled without publishing completion or failure events. Late non-terminal changes (e.g. WAITING from a listen task) are still ignored. Signed-off-by: Eric Deandrea <eric@ericdeandrea.dev> * [Fix #1749] Alternative approach Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> --------- Signed-off-by: Eric Deandrea <eric@ericdeandrea.dev> Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> Co-authored-by: Francisco Javier Tirado Sarti <ftirados@ibm.com>
1 parent 663f046 commit c8f334a

3 files changed

Lines changed: 515 additions & 72 deletions

File tree

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,15 @@ public class LifecycleEventsUtils {
2828

2929
private LifecycleEventsUtils() {}
3030

31+
public static CompletableFuture<Boolean> publishEvent(
32+
boolean statusChanged,
33+
WorkflowContext workflowContext,
34+
Function<WorkflowExecutionCompletableListener, CompletableFuture<?>> function) {
35+
return statusChanged
36+
? publishEvent(workflowContext, function).thenApply(__ -> true)
37+
: CompletableFuture.completedFuture(false);
38+
}
39+
3140
public static CompletableFuture<?> publishEvent(
3241
WorkflowContext workflowContext,
3342
Function<WorkflowExecutionCompletableListener, CompletableFuture<?>> function) {

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

Lines changed: 106 additions & 72 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
import java.util.concurrent.atomic.AtomicReference;
4040
import java.util.concurrent.locks.Lock;
4141
import java.util.concurrent.locks.ReentrantLock;
42+
import java.util.function.Function;
4243
import java.util.function.Supplier;
4344

4445
public class WorkflowMutableInstance implements WorkflowInstance {
@@ -60,7 +61,11 @@ public class WorkflowMutableInstance implements WorkflowInstance {
6061
private Lock statusLock = new ReentrantLock();
6162
private Map<CompletableFuture<TaskContext>, TaskContext> suspended;
6263

63-
private Collection<CompletableFuture<?>> cancelables = new ArrayList<>();
64+
private Collection<CompletableFuture<?>> cancelables =
65+
Collections.synchronizedList(new ArrayList<>());
66+
67+
private Collection<CompletableFuture<Boolean>> outOfOrderListeners =
68+
Collections.synchronizedList(new ArrayList<>());
6469

6570
protected WorkflowMutableInstance(WorkflowDefinition definition, String id, WorkflowModel input) {
6671
this.id = id;
@@ -104,15 +109,16 @@ protected final CompletableFuture<WorkflowModel> startExecution(
104109
.orElse(input))
105110
.whenComplete(this::setCompleteDate)
106111
.thenApply(this::filterAndValidate)
107-
.thenCompose(this::publishEvents)
112+
.thenCompose(this::publishCompletionEvents)
108113
.exceptionallyCompose(this::handleException))
109-
.whenComplete(this::cleanUp);
114+
.handle(this::cleanUpWaitingListeners)
115+
.thenCompose(Function.identity());
110116
futureRef.set(future);
111117
}
112118
return future;
113119
}
114120

115-
private CompletableFuture<WorkflowModel> publishEvents(WorkflowModel model) {
121+
private CompletableFuture<WorkflowModel> publishCompletionEvents(WorkflowModel model) {
116122
return status(WorkflowStatus.COMPLETED)
117123
.thenCompose(
118124
__ ->
@@ -126,19 +132,37 @@ private void setCompleteDate(WorkflowModel result, Throwable ex) {
126132
completedAt = Instant.now();
127133
}
128134

129-
private void cleanUp(WorkflowModel result, Throwable ex) {
135+
private CompletableFuture<WorkflowModel> cleanUpWaitingListeners(
136+
WorkflowModel model, Throwable ex) {
137+
return CompletableFuture.allOf(
138+
outOfOrderListeners.toArray(new CompletableFuture[outOfOrderListeners.size()]))
139+
.handle(this::cleanUp)
140+
.thenCompose(__ -> fromResult(model, ex));
141+
}
142+
143+
private Object cleanUp(Object ignored, Throwable ex) {
130144
additionalObjects.values().stream()
131145
.filter(AutoCloseable.class::isInstance)
132146
.map(AutoCloseable.class::cast)
133147
.forEach(WorkflowUtils::safeClose);
134148
additionalObjects.clear();
135149
workflowContext.definition().removeInstance(this);
150+
outOfOrderListeners.clear();
151+
return ignored;
152+
}
153+
154+
private CompletableFuture<WorkflowModel> fromResult(WorkflowModel model, Throwable ex) {
155+
return ex == null
156+
? CompletableFuture.completedFuture(model)
157+
: CompletableFuture.failedFuture(ex);
136158
}
137159

138160
private CompletableFuture<WorkflowModel> handleException(Throwable exception) {
139161
final Throwable cause =
140162
exception instanceof CompletionException ? exception.getCause() : exception;
141-
if (!(cause instanceof CancellationException)) {
163+
if (cause instanceof CancellationException) {
164+
return CompletableFuture.failedFuture(cause);
165+
} else {
142166
return status(WorkflowStatus.FAULTED)
143167
.thenCompose(
144168
__ ->
@@ -147,7 +171,6 @@ private CompletableFuture<WorkflowModel> handleException(Throwable exception) {
147171
l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, cause))))
148172
.thenCompose(__ -> CompletableFuture.failedFuture(exception));
149173
}
150-
return CompletableFuture.failedFuture(exception);
151174
}
152175

153176
private WorkflowModel filterAndValidate(WorkflowModel model) {
@@ -213,9 +236,21 @@ public <T> T outputAs(Class<T> clazz) {
213236
: null;
214237
}
215238

216-
public CompletableFuture<Boolean> status(WorkflowStatus state) {
217-
WorkflowStatus prevState = this.status.getAndSet(state);
218-
return publishStatusChange(prevState, state);
239+
public CompletableFuture<Boolean> status(WorkflowStatus newState) {
240+
WorkflowStatus prevState;
241+
statusLock.lock();
242+
try {
243+
prevState = status.get();
244+
if (prevState == WorkflowStatus.CANCELLED && newState != WorkflowStatus.CANCELLED) {
245+
return CompletableFuture.failedFuture(
246+
new CancellationException("Workflow has been cancelled"));
247+
} else {
248+
status.set(newState);
249+
}
250+
} finally {
251+
statusLock.unlock();
252+
}
253+
return publishStatusChange(prevState, newState);
219254
}
220255

221256
protected final void setStatus(WorkflowStatus state) {
@@ -224,14 +259,10 @@ protected final void setStatus(WorkflowStatus state) {
224259

225260
private CompletableFuture<Boolean> publishStatusChange(
226261
WorkflowStatus prevState, WorkflowStatus state) {
227-
return prevState != state
228-
? publishEvent(
229-
workflowContext,
230-
l ->
231-
l.onWorkflowStatusChanged(
232-
new WorkflowStatusEvent(workflowContext, prevState, state)))
233-
.thenApply(__ -> true)
234-
: CompletableFuture.completedFuture(false);
262+
return publishEvent(
263+
prevState != state,
264+
workflowContext,
265+
l -> l.onWorkflowStatusChanged(new WorkflowStatusEvent(workflowContext, prevState, state)));
235266
}
236267

237268
@Override
@@ -252,28 +283,30 @@ public boolean suspend() {
252283
WorkflowStatus prevState = internalSuspend();
253284
boolean result = prevState != WorkflowStatus.SUSPENDED;
254285
if (result) {
255-
publishStatusChange(prevState, WorkflowStatus.SUSPENDED)
256-
.thenCompose(
257-
__ ->
258-
publishEvent(
259-
workflowContext,
260-
l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))));
286+
outOfOrder(
287+
publishStatusChange(prevState, WorkflowStatus.SUSPENDED)
288+
.thenCompose(
289+
changed ->
290+
publishEvent(
291+
changed,
292+
workflowContext,
293+
l ->
294+
l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)))));
261295
}
262296
return result;
263297
}
264298

265299
@Override
266300
public CompletableFuture<Boolean> suspendFuture() {
267301
WorkflowStatus prevState = internalSuspend();
268-
return prevState != WorkflowStatus.SUSPENDED
269-
? publishStatusChange(prevState, WorkflowStatus.SUSPENDED)
302+
return outOfOrder(
303+
publishStatusChange(prevState, WorkflowStatus.SUSPENDED)
270304
.thenCompose(
271-
__ ->
305+
changed ->
272306
publishEvent(
307+
changed,
273308
workflowContext,
274-
l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))))
275-
.thenApply(__ -> true)
276-
: CompletableFuture.completedFuture(false);
309+
l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)))));
277310
}
278311

279312
private WorkflowStatus internalSuspend() {
@@ -299,28 +332,29 @@ public boolean resume() {
299332
WorkflowStatus prevStatus = internalResume();
300333
boolean result = prevStatus != WorkflowStatus.RUNNING;
301334
if (result) {
302-
publishStatusChange(prevStatus, WorkflowStatus.RUNNING)
303-
.thenCompose(
304-
__ ->
305-
publishEvent(
306-
workflowContext,
307-
l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))));
335+
outOfOrder(
336+
publishStatusChange(prevStatus, WorkflowStatus.RUNNING)
337+
.thenCompose(
338+
changed ->
339+
publishEvent(
340+
changed,
341+
workflowContext,
342+
l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)))));
308343
}
309344
return result;
310345
}
311346

312347
@Override
313348
public CompletableFuture<Boolean> resumeFuture() {
314349
WorkflowStatus prevStatus = internalResume();
315-
return prevStatus != WorkflowStatus.RUNNING
316-
? publishStatusChange(prevStatus, WorkflowStatus.RUNNING)
350+
return outOfOrder(
351+
publishStatusChange(prevStatus, WorkflowStatus.RUNNING)
317352
.thenCompose(
318-
__ ->
353+
change ->
319354
publishEvent(
355+
change,
320356
workflowContext,
321-
l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))))
322-
.thenApply(__ -> true)
323-
: CompletableFuture.completedFuture(false);
357+
l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)))));
324358
}
325359

326360
private WorkflowStatus internalResume() {
@@ -343,13 +377,12 @@ private WorkflowStatus internalResume() {
343377
return result;
344378
}
345379

346-
public CompletableFuture<TaskContext> cancelCheck(TaskContext t) {
380+
public <T> CompletableFuture<T> cancelCheck(T t) {
347381
try {
348382
statusLock.lock();
349383
if (status.get() == WorkflowStatus.CANCELLED) {
350-
CompletableFuture<TaskContext> cancelled = new CompletableFuture<TaskContext>();
351-
cancelled.completeExceptionally(
352-
new CancellationException("Task " + t.taskName() + " has been cancelled"));
384+
CompletableFuture<T> cancelled = new CompletableFuture<>();
385+
cancelled.completeExceptionally(new CancellationException(t + " has been cancelled"));
353386
return cancelled;
354387
}
355388
} finally {
@@ -380,54 +413,55 @@ public CompletableFuture<TaskContext> suspendedCheck(TaskContext t) {
380413

381414
@Override
382415
public boolean cancel() {
383-
WorkflowStatus prevStatus = internalCancel();
384-
boolean result = prevStatus != WorkflowStatus.CANCELLED;
385-
if (result) {
386-
publishStatusChange(prevStatus, WorkflowStatus.CANCELLED)
387-
.thenCompose(
388-
__ ->
389-
publishEvent(
390-
workflowContext,
391-
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))));
392-
}
393-
return result;
416+
CompletableFuture<Boolean> result = internalCancel();
417+
return !result.isDone() || result.join();
394418
}
395419

396420
@Override
397421
public CompletableFuture<Boolean> cancelFuture() {
398-
WorkflowStatus prevState = internalCancel();
399-
return prevState != WorkflowStatus.CANCELLED
400-
? publishStatusChange(prevState, WorkflowStatus.CANCELLED)
401-
.thenCompose(
402-
__ ->
403-
publishEvent(
404-
workflowContext,
405-
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))))
406-
.thenApply(__ -> true)
407-
: CompletableFuture.completedFuture(false);
422+
return internalCancel();
408423
}
409424

410-
private WorkflowStatus internalCancel() {
411-
WorkflowStatus result;
425+
private CompletableFuture<Boolean> internalCancel() {
426+
WorkflowStatus prevState;
412427
Collection<CompletableFuture<?>> toCancel = null;
413428
try {
414429
statusLock.lock();
415430
if (TaskExecutorHelper.isActive(status.get())) {
416431
toCancel = new ArrayList<>(cancelables);
417432
cancelables.clear();
418-
result = status.getAndSet(WorkflowStatus.CANCELLED);
433+
prevState = status.getAndSet(WorkflowStatus.CANCELLED);
419434
} else {
420-
result = WorkflowStatus.CANCELLED;
435+
prevState = WorkflowStatus.CANCELLED;
421436
}
422437
} finally {
423438
statusLock.unlock();
424439
}
425-
if (result != WorkflowStatus.CANCELLED && toCancel != null) {
440+
CompletableFuture<Boolean> result =
441+
outOfOrder(
442+
publishStatusChange(prevState, WorkflowStatus.CANCELLED)
443+
.thenCompose(
444+
changed ->
445+
publishEvent(
446+
changed,
447+
workflowContext,
448+
l ->
449+
l.onWorkflowCancelled(
450+
new WorkflowCancelledEvent(workflowContext)))));
451+
452+
if (prevState != WorkflowStatus.CANCELLED && toCancel != null) {
426453
toCancel.forEach(t -> t.cancel(true));
427454
}
428455
return result;
429456
}
430457

458+
private CompletableFuture<Boolean> outOfOrder(CompletableFuture<Boolean> future) {
459+
if (!future.isDone()) {
460+
outOfOrderListeners.add(future);
461+
}
462+
return future;
463+
}
464+
431465
public void addCancelable(CompletableFuture<?> cancelable) {
432466
statusLock.lock();
433467
if (status.get() == WorkflowStatus.CANCELLED) {

0 commit comments

Comments
 (0)