@@ -63,7 +63,7 @@ public class WorkflowMutableInstance implements WorkflowInstance {
6363
6464 private Collection <CompletableFuture <?>> cancelables = new ArrayList <>();
6565
66- private volatile CompletableFuture <?> cancelPublished = CompletableFuture . completedFuture ( null );
66+ private Collection < CompletableFuture <Boolean >> outOfOrderListeners = new ArrayList <>( );
6767
6868 protected WorkflowMutableInstance (WorkflowDefinition definition , String id , WorkflowModel input ) {
6969 this .id = id ;
@@ -107,31 +107,16 @@ protected final CompletableFuture<WorkflowModel> startExecution(
107107 .orElse (input ))
108108 .whenComplete (this ::setCompleteDate )
109109 .thenApply (this ::filterAndValidate )
110- .thenCompose (this ::publishEvents )
110+ .thenCompose (this ::publishCompletionEvents )
111111 .exceptionallyCompose (this ::handleException ))
112- .handle (this ::cleanUpAfterCancelPublished )
112+ .handle (this ::cleanUpWaitingListeners )
113113 .thenCompose (Function .identity ());
114114 futureRef .set (future );
115115 }
116116 return future ;
117117 }
118118
119- private CompletableFuture <WorkflowModel > cleanUpAfterCancelPublished (
120- WorkflowModel result , Throwable ex ) {
121- return cancelPublished
122- .handle (
123- (__ , ___ ) -> {
124- cleanUp (result , ex );
125- return null ;
126- })
127- .thenCompose (
128- __ ->
129- ex == null
130- ? CompletableFuture .completedFuture (result )
131- : CompletableFuture .failedFuture (ex ));
132- }
133-
134- private CompletableFuture <WorkflowModel > publishEvents (WorkflowModel model ) {
119+ private CompletableFuture <WorkflowModel > publishCompletionEvents (WorkflowModel model ) {
135120 return status (WorkflowStatus .COMPLETED )
136121 .thenCompose (
137122 __ ->
@@ -145,28 +130,42 @@ private void setCompleteDate(WorkflowModel result, Throwable ex) {
145130 completedAt = Instant .now ();
146131 }
147132
148- private void cleanUp (WorkflowModel result , Throwable ex ) {
133+ private CompletableFuture <WorkflowModel > cleanUpWaitingListeners (
134+ WorkflowModel model , Throwable ex ) {
135+ return CompletableFuture .allOf (
136+ outOfOrderListeners .toArray (new CompletableFuture [outOfOrderListeners .size ()]))
137+ .whenComplete ((__ , ___ ) -> cleanUp ())
138+ .thenCompose (
139+ __ ->
140+ ex == null
141+ ? CompletableFuture .completedFuture (model )
142+ : CompletableFuture .failedFuture (ex ));
143+ }
144+
145+ private void cleanUp () {
149146 additionalObjects .values ().stream ()
150147 .filter (AutoCloseable .class ::isInstance )
151148 .map (AutoCloseable .class ::cast )
152149 .forEach (WorkflowUtils ::safeClose );
153150 additionalObjects .clear ();
154151 workflowContext .definition ().removeInstance (this );
152+ outOfOrderListeners .clear ();
155153 }
156154
157155 private CompletableFuture <WorkflowModel > handleException (Throwable exception ) {
158156 final Throwable cause =
159157 exception instanceof CompletionException ? exception .getCause () : exception ;
160- if (!(cause instanceof CancellationException )) {
158+ if (cause instanceof CancellationException ) {
159+ return CompletableFuture .failedFuture (cause );
160+ } else {
161161 return status (WorkflowStatus .FAULTED )
162162 .thenCompose (
163163 __ ->
164164 publishEvent (
165165 workflowContext ,
166166 l -> l .onWorkflowFailed (new WorkflowFailedEvent (workflowContext , cause ))))
167- .thenCompose (__ -> CompletableFuture .failedFuture (exception ));
167+ .thenCompose (failed -> CompletableFuture .failedFuture (exception ));
168168 }
169- return CompletableFuture .failedFuture (exception );
170169 }
171170
172171 private WorkflowModel filterAndValidate (WorkflowModel model ) {
@@ -232,29 +231,21 @@ public <T> T outputAs(Class<T> clazz) {
232231 : null ;
233232 }
234233
235- public CompletableFuture <Boolean > status (WorkflowStatus state ) {
234+ public CompletableFuture <Boolean > status (WorkflowStatus newState ) {
236235 WorkflowStatus prevState ;
236+ statusLock .lock ();
237237 try {
238- statusLock .lock ();
239- prevState = this .status .get ();
240- if (prevState == WorkflowStatus .CANCELLED ) {
241- // cancellation is final: late callbacks (e.g. an event received by a listen task) must
242- // not move the instance out of it, and a terminal transition must fail so that neither
243- // onWorkflowCompleted nor onWorkflowFailed is published for a cancelled instance
244- return isTerminal (state )
245- ? CompletableFuture .failedFuture (
246- new CancellationException ("Workflow instance " + id + " has been cancelled" ))
247- : CompletableFuture .completedFuture (false );
238+ prevState = status .get ();
239+ if (prevState == WorkflowStatus .CANCELLED && newState != WorkflowStatus .CANCELLED ) {
240+ return CompletableFuture .failedFuture (
241+ new CancellationException ("Workflow has been cancelled" ));
242+ } else {
243+ status .set (newState );
248244 }
249- this .status .set (state );
250245 } finally {
251246 statusLock .unlock ();
252247 }
253- return publishStatusChange (prevState , state );
254- }
255-
256- private static boolean isTerminal (WorkflowStatus state ) {
257- return state == WorkflowStatus .COMPLETED || state == WorkflowStatus .FAULTED ;
248+ return publishStatusChange (prevState , newState );
258249 }
259250
260251 protected final void setStatus (WorkflowStatus state ) {
@@ -263,14 +254,10 @@ protected final void setStatus(WorkflowStatus state) {
263254
264255 private CompletableFuture <Boolean > publishStatusChange (
265256 WorkflowStatus prevState , WorkflowStatus state ) {
266- return prevState != state
267- ? publishEvent (
268- workflowContext ,
269- l ->
270- l .onWorkflowStatusChanged (
271- new WorkflowStatusEvent (workflowContext , prevState , state )))
272- .thenApply (__ -> true )
273- : CompletableFuture .completedFuture (false );
257+ return publishEvent (
258+ prevState != state ,
259+ workflowContext ,
260+ l -> l .onWorkflowStatusChanged (new WorkflowStatusEvent (workflowContext , prevState , state )));
274261 }
275262
276263 @ Override
@@ -291,28 +278,30 @@ public boolean suspend() {
291278 WorkflowStatus prevState = internalSuspend ();
292279 boolean result = prevState != WorkflowStatus .SUSPENDED ;
293280 if (result ) {
294- publishStatusChange (prevState , WorkflowStatus .SUSPENDED )
295- .thenCompose (
296- __ ->
297- publishEvent (
298- workflowContext ,
299- l -> l .onWorkflowSuspended (new WorkflowSuspendedEvent (workflowContext ))));
281+ outOfOrder (
282+ publishStatusChange (prevState , WorkflowStatus .SUSPENDED )
283+ .thenCompose (
284+ changed ->
285+ publishEvent (
286+ changed ,
287+ workflowContext ,
288+ l ->
289+ l .onWorkflowSuspended (new WorkflowSuspendedEvent (workflowContext )))));
300290 }
301291 return result ;
302292 }
303293
304294 @ Override
305295 public CompletableFuture <Boolean > suspendFuture () {
306296 WorkflowStatus prevState = internalSuspend ();
307- return prevState != WorkflowStatus . SUSPENDED
308- ? publishStatusChange (prevState , WorkflowStatus .SUSPENDED )
297+ return outOfOrder (
298+ publishStatusChange (prevState , WorkflowStatus .SUSPENDED )
309299 .thenCompose (
310- __ ->
300+ changed ->
311301 publishEvent (
302+ changed ,
312303 workflowContext ,
313- l -> l .onWorkflowSuspended (new WorkflowSuspendedEvent (workflowContext ))))
314- .thenApply (__ -> true )
315- : CompletableFuture .completedFuture (false );
304+ l -> l .onWorkflowSuspended (new WorkflowSuspendedEvent (workflowContext )))));
316305 }
317306
318307 private WorkflowStatus internalSuspend () {
@@ -338,28 +327,29 @@ public boolean resume() {
338327 WorkflowStatus prevStatus = internalResume ();
339328 boolean result = prevStatus != WorkflowStatus .RUNNING ;
340329 if (result ) {
341- publishStatusChange (prevStatus , WorkflowStatus .RUNNING )
342- .thenCompose (
343- __ ->
344- publishEvent (
345- workflowContext ,
346- l -> l .onWorkflowResumed (new WorkflowResumedEvent (workflowContext ))));
330+ outOfOrder (
331+ publishStatusChange (prevStatus , WorkflowStatus .RUNNING )
332+ .thenCompose (
333+ changed ->
334+ publishEvent (
335+ changed ,
336+ workflowContext ,
337+ l -> l .onWorkflowResumed (new WorkflowResumedEvent (workflowContext )))));
347338 }
348339 return result ;
349340 }
350341
351342 @ Override
352343 public CompletableFuture <Boolean > resumeFuture () {
353344 WorkflowStatus prevStatus = internalResume ();
354- return prevStatus != WorkflowStatus . RUNNING
355- ? publishStatusChange (prevStatus , WorkflowStatus .RUNNING )
345+ return outOfOrder (
346+ publishStatusChange (prevStatus , WorkflowStatus .RUNNING )
356347 .thenCompose (
357- __ ->
348+ change ->
358349 publishEvent (
350+ change ,
359351 workflowContext ,
360- l -> l .onWorkflowResumed (new WorkflowResumedEvent (workflowContext ))))
361- .thenApply (__ -> true )
362- : CompletableFuture .completedFuture (false );
352+ l -> l .onWorkflowResumed (new WorkflowResumedEvent (workflowContext )))));
363353 }
364354
365355 private WorkflowStatus internalResume () {
@@ -382,13 +372,12 @@ private WorkflowStatus internalResume() {
382372 return result ;
383373 }
384374
385- public CompletableFuture <TaskContext > cancelCheck (TaskContext t ) {
375+ public < T > CompletableFuture <T > cancelCheck (T t ) {
386376 try {
387377 statusLock .lock ();
388378 if (status .get () == WorkflowStatus .CANCELLED ) {
389- CompletableFuture <TaskContext > cancelled = new CompletableFuture <TaskContext >();
390- cancelled .completeExceptionally (
391- new CancellationException ("Task " + t .taskName () + " has been cancelled" ));
379+ CompletableFuture <T > cancelled = new CompletableFuture <>();
380+ cancelled .completeExceptionally (new CancellationException (t + " has been cancelled" ));
392381 return cancelled ;
393382 }
394383 } finally {
@@ -419,68 +408,53 @@ public CompletableFuture<TaskContext> suspendedCheck(TaskContext t) {
419408
420409 @ Override
421410 public boolean cancel () {
422- CancelRequest request = internalCancel ();
423- if (request .accepted ()) {
424- publishCancelled (request );
425- }
426- return request .accepted ();
411+ CompletableFuture <Boolean > result = internalCancel ();
412+ return !result .isDone () || result .join ();
427413 }
428414
429415 @ Override
430416 public CompletableFuture <Boolean > cancelFuture () {
431- CancelRequest request = internalCancel ();
432- return request .accepted ()
433- ? publishCancelled (request ).thenApply (__ -> true )
434- : CompletableFuture .completedFuture (false );
435- }
436-
437- /**
438- * Publishes the cancellation and only then cancels the pending futures. Cancelling them may
439- * complete the execution pipeline synchronously, and the pipeline may also end on its own while
440- * the cancellation is being published, so {@link #cleanUpAfterCancelPublished} waits for {@link
441- * CancelRequest#published()} before clearing the instance metadata.
442- */
443- private CompletableFuture <?> publishCancelled (CancelRequest request ) {
444- return publishStatusChange (request .prevStatus (), WorkflowStatus .CANCELLED )
445- .thenCompose (
446- __ ->
447- publishEvent (
448- workflowContext ,
449- l -> l .onWorkflowCancelled (new WorkflowCancelledEvent (workflowContext ))))
450- .whenComplete (
451- (__ , ex ) -> {
452- request .published ().complete (null );
453- request .toCancel ().forEach (t -> t .cancel (true ));
454- });
455- }
456-
457- private record CancelRequest (
458- WorkflowStatus prevStatus ,
459- Collection <CompletableFuture <?>> toCancel ,
460- CompletableFuture <Void > published ) {
461- boolean accepted () {
462- return prevStatus != WorkflowStatus .CANCELLED ;
463- }
417+ return internalCancel ();
464418 }
465419
466- private CancelRequest internalCancel () {
420+ private CompletableFuture <Boolean > internalCancel () {
421+ WorkflowStatus prevState ;
422+ Collection <CompletableFuture <?>> toCancel = null ;
467423 try {
468424 statusLock .lock ();
469425 if (TaskExecutorHelper .isActive (status .get ())) {
470- Collection < CompletableFuture <?>> toCancel = new ArrayList <>(cancelables );
426+ toCancel = new ArrayList <>(cancelables );
471427 cancelables .clear ();
472- CompletableFuture <Void > published = new CompletableFuture <>();
473- cancelPublished = published ;
474- return new CancelRequest (status .getAndSet (WorkflowStatus .CANCELLED ), toCancel , published );
428+ prevState = status .getAndSet (WorkflowStatus .CANCELLED );
475429 } else {
476- return new CancelRequest (
477- WorkflowStatus .CANCELLED ,
478- Collections .emptyList (),
479- CompletableFuture .completedFuture (null ));
430+ prevState = WorkflowStatus .CANCELLED ;
480431 }
481432 } finally {
482433 statusLock .unlock ();
483434 }
435+ CompletableFuture <Boolean > result =
436+ outOfOrder (
437+ publishStatusChange (prevState , WorkflowStatus .CANCELLED )
438+ .thenCompose (
439+ changed ->
440+ publishEvent (
441+ changed ,
442+ workflowContext ,
443+ l ->
444+ l .onWorkflowCancelled (
445+ new WorkflowCancelledEvent (workflowContext )))));
446+
447+ if (prevState != WorkflowStatus .CANCELLED && toCancel != null ) {
448+ toCancel .forEach (t -> t .cancel (true ));
449+ }
450+ return result ;
451+ }
452+
453+ private CompletableFuture <Boolean > outOfOrder (CompletableFuture <Boolean > future ) {
454+ if (!future .isDone ()) {
455+ outOfOrderListeners .add (future );
456+ }
457+ return future ;
484458 }
485459
486460 public void addCancelable (CompletableFuture <?> cancelable ) {
0 commit comments