diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowInstance.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowInstance.java index e4d939073..3c7b912c3 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowInstance.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowInstance.java @@ -51,6 +51,18 @@ public interface WorkflowInstance extends WorkflowInstanceData { boolean resume(); + default CompletableFuture suspendFuture() { + return CompletableFuture.completedFuture(suspend()); + } + + default CompletableFuture cancelFuture() { + return CompletableFuture.completedFuture(cancel()); + } + + default CompletableFuture resumeFuture() { + return CompletableFuture.completedFuture(resume()); + } + T addMetadataIfAbsent(String key, Supplier supplier); void removeMetadata(String key); diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java index c5a499e29..9ab6ecf92 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java @@ -234,13 +234,29 @@ public String toString() { @Override public boolean suspend() { + boolean result = _suspend(); + if (result) { + publishEvent( + workflowContext, l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))); + } + return result; + } + + @Override + public CompletableFuture suspendFuture() { + return _suspend() + ? publishEvent( + workflowContext, + l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))) + .thenApply(__ -> true) + : CompletableFuture.completedFuture(false); + } + + private boolean _suspend() { try { statusLock.lock(); if (TaskExecutorHelper.isActive(status.get()) && suspended == null) { internalSuspend(); - publishEvent( - workflowContext, - l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))); return true; } else { return false; @@ -257,11 +273,29 @@ protected final void internalSuspend() { @Override public boolean resume() { + boolean result = _resume(); + if (result) { + publishEvent( + workflowContext, l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))); + } + return result; + } + + @Override + public CompletableFuture resumeFuture() { + return _resume() + ? publishEvent( + workflowContext, + l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))) + .thenApply(__ -> true) + : CompletableFuture.completedFuture(false); + } + + private boolean _resume() { boolean result; try { statusLock.lock(); if (TaskExecutorHelper.isActive(status.get()) && suspended != null) { - suspended.forEach( (k, v) -> { k.complete(v); @@ -274,10 +308,6 @@ public boolean resume() { } finally { statusLock.unlock(); } - if (result) { - publishEvent( - workflowContext, l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))); - } return result; } @@ -314,6 +344,25 @@ public CompletableFuture suspendedCheck(TaskContext t) { @Override public boolean cancel() { + boolean result = _cancel(); + if (result) { + publishEvent( + workflowContext, l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))); + } + return result; + } + + @Override + public CompletableFuture cancelFuture() { + return _cancel() + ? publishEvent( + workflowContext, + l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))) + .thenApply(__ -> true) + : CompletableFuture.completedFuture(false); + } + + private boolean _cancel() { boolean result; Collection> toCancel = null; try { @@ -330,8 +379,6 @@ public boolean cancel() { statusLock.unlock(); } if (result) { - publishEvent( - workflowContext, l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))); if (toCancel != null) { toCancel.forEach(t -> t.cancel(true)); }