Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Function;
import java.util.function.Supplier;

public class WorkflowMutableInstance implements WorkflowInstance {
Expand All @@ -62,6 +63,8 @@ public class WorkflowMutableInstance implements WorkflowInstance {

private Collection<CompletableFuture<?>> cancelables = new ArrayList<>();

private volatile CompletableFuture<?> cancelPublished = CompletableFuture.completedFuture(null);

protected WorkflowMutableInstance(WorkflowDefinition definition, String id, WorkflowModel input) {
this.id = id;
this.input = input;
Expand Down Expand Up @@ -106,12 +109,28 @@ protected final CompletableFuture<WorkflowModel> startExecution(
.thenApply(this::filterAndValidate)
.thenCompose(this::publishEvents)
.exceptionallyCompose(this::handleException))
.whenComplete(this::cleanUp);
.handle(this::cleanUpAfterCancelPublished)
.thenCompose(Function.identity());
futureRef.set(future);
}
return future;
}

private CompletableFuture<WorkflowModel> cleanUpAfterCancelPublished(
WorkflowModel result, Throwable ex) {
return cancelPublished
.handle(
(__, ___) -> {
cleanUp(result, ex);
return null;
})
.thenCompose(
__ ->
ex == null
? CompletableFuture.completedFuture(result)
: CompletableFuture.failedFuture(ex));
}

private CompletableFuture<WorkflowModel> publishEvents(WorkflowModel model) {
return status(WorkflowStatus.COMPLETED)
.thenCompose(
Expand Down Expand Up @@ -214,10 +233,30 @@ public <T> T outputAs(Class<T> clazz) {
}

public CompletableFuture<Boolean> status(WorkflowStatus state) {
WorkflowStatus prevState = this.status.getAndSet(state);
WorkflowStatus prevState;
try {
statusLock.lock();
prevState = this.status.get();
if (prevState == WorkflowStatus.CANCELLED) {
// cancellation is final: late callbacks (e.g. an event received by a listen task) must
// not move the instance out of it, and a terminal transition must fail so that neither
// onWorkflowCompleted nor onWorkflowFailed is published for a cancelled instance
return isTerminal(state)
? CompletableFuture.failedFuture(
new CancellationException("Workflow instance " + id + " has been cancelled"))
: CompletableFuture.completedFuture(false);
}
this.status.set(state);
} finally {
statusLock.unlock();
}
return publishStatusChange(prevState, state);
}

private static boolean isTerminal(WorkflowStatus state) {
return state == WorkflowStatus.COMPLETED || state == WorkflowStatus.FAULTED;
}

protected final void setStatus(WorkflowStatus state) {
this.status.set(state);
}
Expand Down Expand Up @@ -380,52 +419,68 @@ public CompletableFuture<TaskContext> suspendedCheck(TaskContext t) {

@Override
public boolean cancel() {
WorkflowStatus prevStatus = internalCancel();
boolean result = prevStatus != WorkflowStatus.CANCELLED;
if (result) {
publishStatusChange(prevStatus, WorkflowStatus.CANCELLED)
.thenCompose(
__ ->
publishEvent(
workflowContext,
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))));
CancelRequest request = internalCancel();
if (request.accepted()) {
publishCancelled(request);
}
return result;
return request.accepted();
}

@Override
public CompletableFuture<Boolean> cancelFuture() {
WorkflowStatus prevState = internalCancel();
return prevState != WorkflowStatus.CANCELLED
? publishStatusChange(prevState, WorkflowStatus.CANCELLED)
.thenCompose(
__ ->
publishEvent(
workflowContext,
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))))
.thenApply(__ -> true)
CancelRequest request = internalCancel();
return request.accepted()
? publishCancelled(request).thenApply(__ -> true)
: CompletableFuture.completedFuture(false);
}

private WorkflowStatus internalCancel() {
WorkflowStatus result;
Collection<CompletableFuture<?>> toCancel = null;
/**
* Publishes the cancellation and only then cancels the pending futures. Cancelling them may
* complete the execution pipeline synchronously, and the pipeline may also end on its own while
* the cancellation is being published, so {@link #cleanUpAfterCancelPublished} waits for {@link
* CancelRequest#published()} before clearing the instance metadata.
*/
private CompletableFuture<?> publishCancelled(CancelRequest request) {
return publishStatusChange(request.prevStatus(), WorkflowStatus.CANCELLED)
.thenCompose(
Comment thread
edeandrea marked this conversation as resolved.
Outdated
__ ->
publishEvent(
workflowContext,
l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))))
.whenComplete(
(__, ex) -> {
request.published().complete(null);
request.toCancel().forEach(t -> t.cancel(true));
});
}

private record CancelRequest(
WorkflowStatus prevStatus,
Collection<CompletableFuture<?>> toCancel,
CompletableFuture<Void> published) {
boolean accepted() {
return prevStatus != WorkflowStatus.CANCELLED;
}
}

private CancelRequest internalCancel() {
try {
statusLock.lock();
if (TaskExecutorHelper.isActive(status.get())) {
toCancel = new ArrayList<>(cancelables);
Collection<CompletableFuture<?>> toCancel = new ArrayList<>(cancelables);
cancelables.clear();
result = status.getAndSet(WorkflowStatus.CANCELLED);
CompletableFuture<Void> published = new CompletableFuture<>();
cancelPublished = published;
return new CancelRequest(status.getAndSet(WorkflowStatus.CANCELLED), toCancel, published);
} else {
result = WorkflowStatus.CANCELLED;
return new CancelRequest(
WorkflowStatus.CANCELLED,
Collections.emptyList(),
CompletableFuture.completedFuture(null));
}
} finally {
statusLock.unlock();
}
if (result != WorkflowStatus.CANCELLED && toCancel != null) {
toCancel.forEach(t -> t.cancel(true));
}
return result;
}

public void addCancelable(CompletableFuture<?> cancelable) {
Expand Down
Loading
Loading