diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/AbstractExecutorServiceHolder.java b/impl/core/src/main/java/io/serverlessworkflow/impl/AbstractExecutorServiceHolder.java deleted file mode 100644 index 9982c9f29..000000000 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/AbstractExecutorServiceHolder.java +++ /dev/null @@ -1,32 +0,0 @@ -/* - * Copyright 2020-Present The Serverless Workflow Specification Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.serverlessworkflow.impl; - -import java.util.concurrent.ExecutorService; -import java.util.concurrent.TimeUnit; - -public abstract class AbstractExecutorServiceHolder implements ExecutorServiceFactory { - - protected ExecutorService service; - - @Override - public void close() throws InterruptedException { - if (service != null && !service.isShutdown()) { - service.shutdown(); - service.awaitTermination(2, TimeUnit.SECONDS); - } - } -} diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java b/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java index 7710c43fc..8d43b7d7d 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/DefaultExecutorServiceFactory.java @@ -17,23 +17,21 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; +import java.util.concurrent.TimeUnit; -public class DefaultExecutorServiceFactory extends AbstractExecutorServiceHolder { - - private Lock serviceLock = new ReentrantLock(); +public class DefaultExecutorServiceFactory implements ExecutorServiceFactory { + private ExecutorService service = Executors.newCachedThreadPool(); @Override public ExecutorService get() { - try { - serviceLock.lock(); - if (service == null) { - service = Executors.newCachedThreadPool(); - } - } finally { - serviceLock.unlock(); - } return service; } + + @Override + public void close() throws Exception { + if (!service.isShutdown()) { + service.shutdown(); + service.awaitTermination(2, TimeUnit.SECONDS); + } + } } diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/ExecutorServiceHolder.java b/impl/core/src/main/java/io/serverlessworkflow/impl/ExecutorServiceHolder.java deleted file mode 100644 index 23eb9c691..000000000 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/ExecutorServiceHolder.java +++ /dev/null @@ -1,30 +0,0 @@ -/* - * Copyright 2020-Present The Serverless Workflow Specification Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.serverlessworkflow.impl; - -import java.util.concurrent.ExecutorService; - -public class ExecutorServiceHolder extends AbstractExecutorServiceHolder { - - public ExecutorServiceHolder(ExecutorService service) { - this.service = service; - } - - @Override - public ExecutorService get() { - return service; - } -} diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java index e393dcda0..59b3dbe4b 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java @@ -254,7 +254,7 @@ public SchemaValidator getValidator(SchemaInline inline) { private WorkflowPositionFactory positionFactory = () -> new QueueWorkflowPosition(); private WorkflowInstanceIdFactory idFactory; private WorkflowScheduler scheduler; - private ExecutorServiceFactory executorFactory = new DefaultExecutorServiceFactory(); + private ExecutorServiceFactory executorFactory; private EventConsumer eventConsumer; private Collection eventPublishers = new ArrayList<>(); private RuntimeDescriptorFactory descriptorFactory = @@ -458,7 +458,9 @@ public Builder withCronResolverFactory(CronResolverFactory cronResolverFactory) } public WorkflowApplication build() { - + if (executorFactory == null) { + executorFactory = new DefaultExecutorServiceFactory(); + } if (modelFactory == null) { modelFactory = loadFirst(WorkflowModelFactory.class) 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 9ab6ecf92..fa0b5ba4a 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java @@ -74,49 +74,53 @@ public CompletableFuture start() { return startExecution( () -> { startedAt = Instant.now(); - return publishEvent( - workflowContext, l -> l.onWorkflowStarted(new WorkflowStartedEvent(workflowContext))); + return status(WorkflowStatus.RUNNING) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowStarted(new WorkflowStartedEvent(workflowContext)))); }); } protected final CompletableFuture startExecution( Supplier> runnable) { CompletableFuture future = futureRef.get(); - if (future != null) { - return future; + if (future == null) { + future = + runnable + .get() + .thenCompose( + v -> + TaskExecutorHelper.processTaskList( + workflowContext.definition().startTask(), + workflowContext, + Optional.empty(), + workflowContext + .definition() + .inputFilter() + .map(f -> f.apply(workflowContext, null, input)) + .orElse(input)) + .whenComplete(this::setCompleteDate) + .thenApply(this::filterAndValidate) + .thenCompose(this::publishEvents) + .exceptionallyCompose(this::handleException)) + .whenComplete(this::cleanUp); + futureRef.set(future); } - status(WorkflowStatus.RUNNING); - - future = - runnable - .get() - .thenCompose( - v -> - TaskExecutorHelper.processTaskList( - workflowContext.definition().startTask(), - workflowContext, - Optional.empty(), - workflowContext - .definition() - .inputFilter() - .map(f -> f.apply(workflowContext, null, input)) - .orElse(input)) - .whenComplete(this::setCompleteDate) - .thenApply(this::filterAndValidate) - .whenComplete(this::handleException) - .thenCompose( - model -> - publishEvent( - workflowContext, - l -> - l.onWorkflowCompleted( - new WorkflowCompletedEvent(workflowContext, model))) - .thenApply(__ -> model))) - .whenComplete(this::cleanUp); - futureRef.set(future); return future; } + private CompletableFuture publishEvents(WorkflowModel model) { + return status(WorkflowStatus.COMPLETED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowCompleted(new WorkflowCompletedEvent(workflowContext, model)))) + .thenApply(__ -> model); + } + private void setCompleteDate(WorkflowModel result, Throwable ex) { completedAt = Instant.now(); } @@ -130,19 +134,19 @@ private void cleanUp(WorkflowModel result, Throwable ex) { workflowContext.definition().removeInstance(this); } - private void handleException(WorkflowModel result, Throwable exception) { - if (exception != null) { - final Throwable cause = - exception instanceof CompletionException ? exception.getCause() : exception; - if (!(cause instanceof CancellationException)) { - status(WorkflowStatus.FAULTED); - publishEvent( - workflowContext, - l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, cause))); - } - } else { - status(WorkflowStatus.COMPLETED); + private CompletableFuture handleException(Throwable exception) { + final Throwable cause = + exception instanceof CompletionException ? exception.getCause() : exception; + if (!(cause instanceof CancellationException)) { + return status(WorkflowStatus.FAULTED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, cause)))) + .thenCompose(__ -> CompletableFuture.failedFuture(exception)); } + return CompletableFuture.failedFuture(exception); } private WorkflowModel filterAndValidate(WorkflowModel model) { @@ -208,15 +212,25 @@ public T outputAs(Class clazz) { : null; } - public void status(WorkflowStatus state) { + public CompletableFuture status(WorkflowStatus state) { WorkflowStatus prevState = this.status.getAndSet(state); - if (prevState != state) { - publishEvent( - workflowContext, - l -> - l.onWorkflowStatusChanged( - new WorkflowStatusEvent(workflowContext, prevState, state))); - } + return publishStatusChange(prevState, state); + } + + protected final void setStatus(WorkflowStatus state) { + this.status.set(state); + } + + private CompletableFuture publishStatusChange( + WorkflowStatus prevState, WorkflowStatus state) { + return prevState != state + ? publishEvent( + workflowContext, + l -> + l.onWorkflowStatusChanged( + new WorkflowStatusEvent(workflowContext, prevState, state))) + .thenApply(__ -> true) + : CompletableFuture.completedFuture(false); } @Override @@ -234,65 +248,82 @@ public String toString() { @Override public boolean suspend() { - boolean result = _suspend(); + WorkflowStatus prevState = internalSuspend(); + boolean result = prevState != WorkflowStatus.SUSPENDED; if (result) { - publishEvent( - workflowContext, l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))); + publishStatusChange(prevState, WorkflowStatus.SUSPENDED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)))); } return result; } @Override public CompletableFuture suspendFuture() { - return _suspend() - ? publishEvent( - workflowContext, - l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext))) + WorkflowStatus prevState = internalSuspend(); + return prevState != WorkflowStatus.SUSPENDED + ? publishStatusChange(prevState, WorkflowStatus.SUSPENDED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowSuspended(new WorkflowSuspendedEvent(workflowContext)))) .thenApply(__ -> true) : CompletableFuture.completedFuture(false); } - private boolean _suspend() { + private WorkflowStatus internalSuspend() { try { statusLock.lock(); if (TaskExecutorHelper.isActive(status.get()) && suspended == null) { - internalSuspend(); - return true; + setSuspended(); + return status.getAndSet(WorkflowStatus.SUSPENDED); } else { - return false; + return WorkflowStatus.SUSPENDED; } } finally { statusLock.unlock(); } } - protected final void internalSuspend() { + protected final void setSuspended() { suspended = new ConcurrentHashMap<>(); - status(WorkflowStatus.SUSPENDED); } @Override public boolean resume() { - boolean result = _resume(); + WorkflowStatus prevStatus = internalResume(); + boolean result = prevStatus != WorkflowStatus.RUNNING; if (result) { - publishEvent( - workflowContext, l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))); + publishStatusChange(prevStatus, WorkflowStatus.RUNNING) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)))); } return result; } @Override public CompletableFuture resumeFuture() { - return _resume() - ? publishEvent( - workflowContext, - l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext))) + WorkflowStatus prevStatus = internalResume(); + return prevStatus != WorkflowStatus.RUNNING + ? publishStatusChange(prevStatus, WorkflowStatus.RUNNING) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowResumed(new WorkflowResumedEvent(workflowContext)))) .thenApply(__ -> true) : CompletableFuture.completedFuture(false); } - private boolean _resume() { - boolean result; + private WorkflowStatus internalResume() { + WorkflowStatus result; try { statusLock.lock(); if (TaskExecutorHelper.isActive(status.get()) && suspended != null) { @@ -301,9 +332,9 @@ private boolean _resume() { k.complete(v); }); suspended = null; - result = true; + result = status.getAndSet(WorkflowStatus.RUNNING); } else { - result = false; + result = WorkflowStatus.RUNNING; } } finally { statusLock.unlock(); @@ -327,61 +358,71 @@ public CompletableFuture cancelCheck(TaskContext t) { } public CompletableFuture suspendedCheck(TaskContext t) { + final WorkflowStatus prevState; try { statusLock.lock(); if (suspended != null) { CompletableFuture suspendedTask = new CompletableFuture(); suspended.put(suspendedTask, t); + prevState = WorkflowStatus.RUNNING; return suspendedTask; } else if (TaskExecutorHelper.isActive(status.get())) { - status(WorkflowStatus.RUNNING); + prevState = this.status.getAndSet(WorkflowStatus.RUNNING); + } else { + prevState = WorkflowStatus.RUNNING; } } finally { statusLock.unlock(); } - return CompletableFuture.completedFuture(t); + return publishStatusChange(prevState, WorkflowStatus.RUNNING).thenApply(__ -> t); } @Override public boolean cancel() { - boolean result = _cancel(); + WorkflowStatus prevStatus = internalCancel(); + boolean result = prevStatus != WorkflowStatus.CANCELLED; if (result) { - publishEvent( - workflowContext, l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))); + publishStatusChange(prevStatus, WorkflowStatus.CANCELLED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext)))); } return result; } @Override public CompletableFuture cancelFuture() { - return _cancel() - ? publishEvent( - workflowContext, - l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext))) + WorkflowStatus prevState = internalCancel(); + return prevState != WorkflowStatus.CANCELLED + ? publishStatusChange(prevState, WorkflowStatus.CANCELLED) + .thenCompose( + __ -> + publishEvent( + workflowContext, + l -> l.onWorkflowCancelled(new WorkflowCancelledEvent(workflowContext)))) .thenApply(__ -> true) : CompletableFuture.completedFuture(false); } - private boolean _cancel() { - boolean result; + private WorkflowStatus internalCancel() { + WorkflowStatus result; Collection> toCancel = null; try { statusLock.lock(); if (TaskExecutorHelper.isActive(status.get())) { toCancel = new ArrayList<>(cancelables); cancelables.clear(); - status(WorkflowStatus.CANCELLED); - result = true; + result = status.getAndSet(WorkflowStatus.CANCELLED); } else { - result = false; + result = WorkflowStatus.CANCELLED; } } finally { statusLock.unlock(); } - if (result) { - if (toCancel != null) { - toCancel.forEach(t -> t.cancel(true)); - } + if (result != WorkflowStatus.CANCELLED && toCancel != null) { + toCancel.forEach(t -> t.cancel(true)); } return result; } diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/AbstractTaskExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/AbstractTaskExecutor.java index af335f7a9..58d1eb5f2 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/AbstractTaskExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/AbstractTaskExecutor.java @@ -193,7 +193,7 @@ private CompletableFuture executeNext( WorkflowContext workflow, TaskContext taskContext) { TransitionInfo transition = taskContext.transition(); if (transition.isEndNode()) { - workflow.instance().status(WorkflowStatus.COMPLETED); + return workflow.instance().status(WorkflowStatus.COMPLETED).thenApply(__ -> taskContext); } else if (transition.next() != null) { return transition.next().apply(workflow, taskContext.parent(), taskContext.output()); } @@ -241,15 +241,6 @@ public CompletableFuture apply( }) .thenCompose(t -> execute(workflowContext, t)) .thenCompose(workflowContext.instance()::cancelCheck) - .whenComplete( - (t, e) -> { - if (e != null) { - handleException( - workflowContext, - taskContext, - e instanceof CompletionException ? e.getCause() : e); - } - }) .thenApply( t -> { outputProcessor.ifPresent( @@ -272,7 +263,13 @@ public CompletableFuture apply( l -> l.onTaskCompleted( new TaskCompletedEvent(workflowContext, taskContext))) - .thenApply(__ -> t)); + .thenApply(__ -> t)) + .exceptionallyCompose( + e -> + handleException( + workflowContext, + taskContext, + e instanceof CompletionException ? e.getCause() : e)); if (timeout.isPresent()) { completable = completable @@ -298,17 +295,21 @@ public CompletableFuture apply( } } - private void handleException( + private CompletableFuture handleException( WorkflowContext workflowContext, TaskContext taskContext, Throwable e) { + CompletableFuture events; if (e instanceof CancellationException) { - publishEvent( - workflowContext, - l -> l.onTaskCancelled(new TaskCancelledEvent(workflowContext, taskContext))); + events = + publishEvent( + workflowContext, + l -> l.onTaskCancelled(new TaskCancelledEvent(workflowContext, taskContext))); } else { - publishEvent( - workflowContext, - l -> l.onTaskFailed(new TaskFailedEvent(workflowContext, taskContext, e))); + events = + publishEvent( + workflowContext, + l -> l.onTaskFailed(new TaskFailedEvent(workflowContext, taskContext, e))); } + return events.thenCompose(__ -> CompletableFuture.failedFuture(e)); } public WorkflowPosition position() { diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ListenExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ListenExecutor.java index 6511a94a7..80d0ed8ac 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ListenExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ListenExecutor.java @@ -35,6 +35,8 @@ import io.serverlessworkflow.impl.events.EventRegistrationBuilderCollection; import io.serverlessworkflow.impl.events.EventRegistrationBuilderInfo; import io.serverlessworkflow.impl.events.EventRegistrationInfo; +import java.util.ArrayList; +import java.util.Collection; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.function.BiConsumer; @@ -96,7 +98,8 @@ protected void internalProcessCe( WorkflowModelCollection arrayNode, WorkflowContext workflow, TaskContext taskContext, - CompletableFuture future) { + CompletableFuture future, + Collection> waitingListeners) { arrayNode.add(node); future.complete(node); } @@ -143,13 +146,15 @@ protected void internalProcessCe( WorkflowModelCollection arrayNode, WorkflowContext workflow, TaskContext taskContext, - CompletableFuture future) { + CompletableFuture future, + Collection> waitingListeners) { arrayNode.add(node); if (until.map(u -> u.test(workflow, taskContext, arrayNode)).orElse(true) && untilRegBuilders == null) { future.complete(node); } else { - ((WorkflowMutableInstance) workflow.instance()).status(WorkflowStatus.WAITING); + waitingListeners.add( + ((WorkflowMutableInstance) workflow.instance()).status(WorkflowStatus.WAITING)); } } } @@ -159,25 +164,38 @@ protected abstract void internalProcessCe( WorkflowModelCollection arrayNode, WorkflowContext workflow, TaskContext taskContext, - CompletableFuture future); + CompletableFuture future, + Collection> waitingListeners); @Override protected CompletableFuture internalExecute( WorkflowContext workflow, TaskContext taskContext) { WorkflowModelCollection output = workflow.definition().application().modelFactory().createCollection(); - ((WorkflowMutableInstance) workflow.instance()).status(WorkflowStatus.WAITING); + Collection> waitingListeners = new ArrayList<>(); + waitingListeners.add( + ((WorkflowMutableInstance) workflow.instance()).status(WorkflowStatus.WAITING)); EventRegistrationInfo info = buildInfo( (BiConsumer>) ((ce, future) -> - processCe(converter.apply(ce), output, workflow, taskContext, future)), + processCe( + converter.apply(ce), + output, + workflow, + taskContext, + future, + waitingListeners)), workflow, taskContext); workflow.instance().addCancelable(info.completableFuture()); return info.completableFuture() .whenComplete((__, e) -> info.registrations().forEach(eventConsumer::unregister)) - .thenApply(__ -> output); + .thenCompose( + __ -> + CompletableFuture.allOf( + waitingListeners.toArray(new CompletableFuture[waitingListeners.size()]))) + .handle((__, ___) -> output); } protected EventRegistrationInfo buildInfo( @@ -193,7 +211,8 @@ private void processCe( WorkflowModelCollection arrayNode, WorkflowContext workflow, TaskContext taskContext, - CompletableFuture future) { + CompletableFuture future, + Collection> waitingListeners) { loop.ifPresentOrElse( t -> { SubscriptionIterator forEach = task.getForeach(); @@ -206,9 +225,12 @@ private void processCe( taskContext.variables().put(at, arrayNode.size()); } TaskExecutorHelper.processTaskList(t, workflow, Optional.of(taskContext), node) - .thenAccept(n -> internalProcessCe(n, arrayNode, workflow, taskContext, future)); + .thenAccept( + n -> + internalProcessCe( + n, arrayNode, workflow, taskContext, future, waitingListeners)); }, - () -> internalProcessCe(node, arrayNode, workflow, taskContext, future)); + () -> internalProcessCe(node, arrayNode, workflow, taskContext, future, waitingListeners)); } protected ListenExecutor(ListenExecutorBuilder builder) { diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/WaitExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/WaitExecutor.java index e6bc9366f..94e0796f8 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/WaitExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/WaitExecutor.java @@ -56,13 +56,14 @@ protected WaitExecutor(WaitExecutorBuilder builder) { @Override protected CompletableFuture internalExecute( WorkflowContext workflow, TaskContext taskContext) { - workflow.instance().status(WorkflowStatus.WAITING); + CompletableFuture listenerFuture = workflow.instance().status(WorkflowStatus.WAITING); CompletableFuture future = new CompletableFuture<>(); CompletableFuture.delayedExecutor( durationResolver.apply(workflow, taskContext, taskContext.input()).toMillis(), TimeUnit.MILLISECONDS, workflow.definition().application().executorService()) - .execute(() -> future.complete(taskContext.output())); + .execute( + () -> listenerFuture.whenComplete((__, ___) -> future.complete(taskContext.output()))); return future; } } diff --git a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java index a54ed9b1f..b132c937d 100644 --- a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java +++ b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java @@ -55,7 +55,10 @@ public CompletableFuture start() { return startExecution( () -> { if (info.status() == WorkflowStatus.SUSPENDED) { - internalSuspend(); + setSuspended(); + setStatus(WorkflowStatus.SUSPENDED); + } else { + setStatus(WorkflowStatus.RUNNING); } return CompletableFuture.completedFuture(null); }); diff --git a/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java b/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java index 892be2950..26210cdc2 100644 --- a/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java +++ b/impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java @@ -25,7 +25,6 @@ import io.serverlessworkflow.api.types.Workflow; import io.serverlessworkflow.fluent.spec.WorkflowBuilder; import io.serverlessworkflow.fluent.spec.dsl.DSL; -import io.serverlessworkflow.impl.ExecutorServiceFactory; import io.serverlessworkflow.impl.WorkflowApplication; import io.serverlessworkflow.impl.WorkflowDefinition; import io.serverlessworkflow.impl.WorkflowDefinitionId; @@ -55,8 +54,6 @@ import java.util.concurrent.CompletionException; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.stream.Stream; @@ -83,21 +80,6 @@ static void init() { appl = WorkflowApplication.builder() .withLifeCycleCloudEventFactory(new InputOutputLifeCycleCloudEventFactory()) - .withExecutorFactory( - new ExecutorServiceFactory() { - - private ExecutorService service = Executors.newFixedThreadPool(2); - - @Override - public void close() throws Exception { - service.shutdownNow(); - } - - @Override - public ExecutorService get() { - return service; - } - }) .withEventConsumer(eventBroker) .withEventPublisher(eventBroker) .build();