Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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 @@ -19,9 +19,13 @@
import io.serverlessworkflow.impl.WorkflowInstance;
import io.serverlessworkflow.impl.WorkflowModel;
import java.util.function.Consumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class ScheduledInstanceRunnable implements Runnable, Consumer<WorkflowInstance> {

private static final Logger logger = LoggerFactory.getLogger(ScheduledInstanceRunnable.class);

protected final WorkflowDefinition definition;

public ScheduledInstanceRunnable(WorkflowDefinition definition) {
Expand All @@ -45,6 +49,22 @@ public static void runScheduledInstance(WorkflowDefinition definition, WorkflowM
private static void runScheduledInstance(
WorkflowDefinition definition, WorkflowInstance instance) {
definition.addScheduledInstance(instance);
definition.application().executorService().execute(() -> instance.start());
definition
.application()
.executorService()
.execute(
() ->
instance
.start()
.whenComplete(
(model, ex) -> {
if (ex != null) {
logger.error(
"Scheduled workflow instance {} of definition {} has failed",
instance.id(),
definition.id(),
ex);
}
}));
Comment on lines +59 to +68
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
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;
Expand Down Expand Up @@ -54,6 +55,8 @@
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;
Expand All @@ -76,9 +79,25 @@ static void init() {
for (String type : AbstractLifeCyclePublisher.getLifeCycleTypes()) {
eventBroker.register(type, ce -> publishedEvents.add(ce));
}
;
appl =
WorkflowApplication.builder()
.withLifeCycleCloudEventFactory(new InputOutputLifeCycleCloudEventFactory())
.withExecutorFactory(
new ExecutorServiceFactory() {

private ExecutorService service = Executors.newFixedThreadPool(2);

@Override
public void close() throws Exception {
service.shutdownNow();
}
Comment on lines +92 to +94

@Override
public ExecutorService get() {
return service;
}
})
.withEventConsumer(eventBroker)
.withEventPublisher(eventBroker)
.build();
Expand Down Expand Up @@ -158,34 +177,6 @@ private void doTestSuspendResumeNotWait(Workflow workflow)
.isBeforeOrEqualTo(workflowResumedEvent.resumedAt());
}

@ParameterizedTest(name = "{0}")
@MethodSource("waitSetWorkflowSources")
void testSuspendResumeWait(String sourceName, Workflow workflow)
throws ExecutionException, InterruptedException, TimeoutException {
doTestSuspendResumeWait(workflow);
}

private void doTestSuspendResumeWait(Workflow workflow)
throws InterruptedException, ExecutionException, TimeoutException {
WorkflowDefinition def = appl.workflowDefinition(workflow);
WorkflowInstance instance = def.instance(Map.of());
CompletableFuture<WorkflowModel> future = instance.start();
assertThat(instance.status()).isEqualTo(WorkflowStatus.WAITING);
instance.suspend();
assertThat(instance.status()).isEqualTo(WorkflowStatus.SUSPENDED);
Thread.sleep(550);
instance.resume();
assertThat(future.get(1, TimeUnit.SECONDS).asMap().orElseThrow())
.isEqualTo(Map.of("name", "Javierito"));
assertThat(instance.status()).isEqualTo(WorkflowStatus.COMPLETED);
WorkflowSuspendedCEData workflowSuspendedEvent =
assertPojoInCE(
"io.serverlessworkflow.workflow.suspended.v1", WorkflowSuspendedCEData.class);
WorkflowResumedCEData workflowResumedEvent =
assertPojoInCE("io.serverlessworkflow.workflow.resumed.v1", WorkflowResumedCEData.class);
assertThat(workflowSuspendedEvent.suspendedAt()).isBefore(workflowResumedEvent.resumedAt());
}

@ParameterizedTest(name = "{0}")
@MethodSource("waitSetWorkflowSources")
void testCancel(String sourceName, Workflow workflow) throws IOException {
Expand All @@ -206,16 +197,19 @@ private void doTestCancel(Workflow workflow) {

@ParameterizedTest(name = "{0}")
@MethodSource("waitSetWorkflowSources")
void testSuspendResumeTimeout(String sourceName, Workflow workflow) {
doTestSuspendResumeTimeout(workflow);
void testSuspendTimeout(String sourceName, Workflow workflow) {
doTestSuspendTimeout(workflow);
}

private static void doTestSuspendResumeTimeout(Workflow workflow) {
private static void doTestSuspendTimeout(Workflow workflow) {
WorkflowDefinition def = appl.workflowDefinition(workflow);
WorkflowInstance instance = def.instance(Map.of());
CompletableFuture<WorkflowModel> future = instance.start();
instance.suspend();
assertThat(catchThrowableOfType(TimeoutException.class, () -> future.get(1, TimeUnit.SECONDS)))
assertThat(instance.status()).isEqualTo(WorkflowStatus.SUSPENDED);
assertThat(
catchThrowableOfType(
TimeoutException.class, () -> future.get(400, TimeUnit.MILLISECONDS)))
.isNotNull();
Comment thread
fjtirado marked this conversation as resolved.
Comment thread
fjtirado marked this conversation as resolved.
}

Expand Down Expand Up @@ -255,11 +249,10 @@ private <T> T assertPojoInCE(String type, Class<T> clazz) {
private static Workflow waitTestWorkflow() {
return WorkflowBuilder.workflow("wait-test-java-dsl", "test", "0.1.0")
.tasks(
// wait 500 ms
DSL.wait(
"waitABit",
timeoutBuilder ->
timeoutBuilder.duration(durationBuilder -> durationBuilder.milliseconds(500))),
timeoutBuilder.duration(durationBuilder -> durationBuilder.milliseconds(200))),
DSL.set("useExpression", setTaskBuilder -> setTaskBuilder.put("name", "Javierito")))
.build();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ document:
do:
- waitABit:
wait:
milliseconds: 500
milliseconds: 200
- useExpression:
set:
name : Javierito
Loading