From e4e9b732a2661a9d4a4200658aec175cddb0762a Mon Sep 17 00:00:00 2001 From: Francisco Javier Tirado Sarti Date: Fri, 7 Aug 2026 14:30:51 +0200 Subject: [PATCH 1/2] [Fix #1601] Implement then on catch Signed-off-by: Francisco Javier Tirado Sarti --- .../impl/executors/CallTaskExecutor.java | 2 +- .../impl/executors/DoExecutor.java | 2 +- .../impl/executors/EmitExecutor.java | 3 +- .../impl/executors/ForExecutor.java | 2 +- .../impl/executors/ForkExecutor.java | 3 +- .../impl/executors/ListenExecutor.java | 3 +- .../impl/executors/RaiseExecutor.java | 3 +- .../impl/executors/RegularTaskExecutor.java | 17 +++++----- .../impl/executors/RunTaskExecutor.java | 3 +- .../impl/executors/SetExecutor.java | 2 +- .../impl/executors/TryExecutor.java | 31 ++++++++++++++++++- .../impl/executors/WaitExecutor.java | 3 +- .../impl/test/RetryTimeoutTest.java | 1 + .../try-catch-match-then.yaml | 24 ++++++++++++++ 14 files changed, 80 insertions(+), 19 deletions(-) create mode 100644 impl/test/src/test/resources/workflows-samples/try-catch-match-then.yaml diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/CallTaskExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/CallTaskExecutor.java index fec6bcb7e..556f12253 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/CallTaskExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/CallTaskExecutor.java @@ -29,7 +29,7 @@ public class CallTaskExecutor extends RegularTaskExecutor private final CallableTask callable; public static class CallTaskExecutorBuilder - extends RegularTaskExecutorBuilder { + extends RegularTaskExecutorBuilder> { private CallableTaskFactory callableFactory; private List callableProxyBuilders; private CallableTask callable; diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/DoExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/DoExecutor.java index c5a1b95f2..4b351b101 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/DoExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/DoExecutor.java @@ -28,7 +28,7 @@ public class DoExecutor extends RegularTaskExecutor { private final TaskExecutor taskExecutor; - public static class DoExecutorBuilder extends RegularTaskExecutorBuilder { + public static class DoExecutorBuilder extends RegularTaskExecutorBuilder { private TaskExecutor taskExecutor; protected DoExecutorBuilder( diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/EmitExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/EmitExecutor.java index 483754847..f34253471 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/EmitExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/EmitExecutor.java @@ -53,7 +53,8 @@ public class EmitExecutor extends RegularTaskExecutor { .toList(); private final EventPropertiesBuilder props; - public static class EmitExecutorBuilder extends RegularTaskExecutorBuilder { + public static class EmitExecutorBuilder + extends RegularTaskExecutorBuilder { private EventPropertiesBuilder eventBuilder; diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForExecutor.java index 6022002e3..4eeb69ff7 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForExecutor.java @@ -37,7 +37,7 @@ public class ForExecutor extends RegularTaskExecutor { private final Optional whileExpr; private final TaskExecutor taskExecutor; - public static class ForExecutorBuilder extends RegularTaskExecutorBuilder { + public static class ForExecutorBuilder extends RegularTaskExecutorBuilder { private TaskExecutor taskExecutor; protected ForExecutorBuilder( diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForkExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForkExecutor.java index 786632ae4..638a53c02 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForkExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/ForkExecutor.java @@ -39,7 +39,8 @@ public class ForkExecutor extends RegularTaskExecutor { private final boolean compete; - public static class ForkExecutorBuilder extends RegularTaskExecutorBuilder { + public static class ForkExecutorBuilder + extends RegularTaskExecutorBuilder { private final Map> taskExecutors; private final boolean compete; 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 5896e8a6f..6511a94a7 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 @@ -47,7 +47,8 @@ public abstract class ListenExecutor extends RegularTaskExecutor { protected final Function converter; protected final EventConsumer eventConsumer; - public static class ListenExecutorBuilder extends RegularTaskExecutorBuilder { + public static class ListenExecutorBuilder + extends RegularTaskExecutorBuilder { private EventRegistrationBuilderInfo registrationInfo; private TaskExecutor loop; diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RaiseExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RaiseExecutor.java index 46b0cfe59..4c179750c 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RaiseExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RaiseExecutor.java @@ -42,7 +42,8 @@ public class RaiseExecutor extends RegularTaskExecutor { private final BiFunction errorBuilder; - public static class RaiseExecutorBuilder extends RegularTaskExecutorBuilder { + public static class RaiseExecutorBuilder + extends RegularTaskExecutorBuilder { private final BiFunction errorBuilder; private final WorkflowValueResolver typeFilter; diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RegularTaskExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RegularTaskExecutor.java index 9213dc3f3..691783040 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RegularTaskExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RegularTaskExecutor.java @@ -28,12 +28,14 @@ public abstract class RegularTaskExecutor extends AbstractTa protected TransitionInfo transition; - protected RegularTaskExecutor(RegularTaskExecutorBuilder builder) { + protected > RegularTaskExecutor( + RegularTaskExecutorBuilder builder) { super(builder); } - public abstract static class RegularTaskExecutorBuilder - extends AbstractTaskExecutorBuilder> { + public abstract static class RegularTaskExecutorBuilder< + T extends TaskBase, V extends RegularTaskExecutor> + extends AbstractTaskExecutorBuilder { private TransitionInfoBuilder transition; @@ -42,12 +44,13 @@ protected RegularTaskExecutorBuilder( super(position, task, definition); } + @Override public void connect(Map> connections) { this.transition = next(task.getThen(), connections); } @Override - protected void buildTransition(RegularTaskExecutor instance) { + protected void buildTransition(V instance) { instance.transition = TransitionInfo.build(transition); } } @@ -59,10 +62,8 @@ protected TransitionInfo getSkipTransition() { protected CompletableFuture execute( WorkflowContext workflow, TaskContext taskContext) { - CompletableFuture future = - internalExecute(workflow, taskContext) - .thenApply(node -> taskContext.rawOutput(node).transition(transition)); - return future; + return internalExecute(workflow, taskContext) + .thenApply(node -> taskContext.rawOutput(node).transition(transition)); } protected abstract CompletableFuture internalExecute( diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RunTaskExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RunTaskExecutor.java index af2a414d5..c398a36d7 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RunTaskExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/RunTaskExecutor.java @@ -33,7 +33,8 @@ public class RunTaskExecutor extends RegularTaskExecutor { private static final ServiceLoader runnables = ServiceLoader.load(RunnableTaskBuilder.class); - public static class RunTaskExecutorBuilder extends RegularTaskExecutorBuilder { + public static class RunTaskExecutorBuilder + extends RegularTaskExecutorBuilder { private CallableTask runnable; protected RunTaskExecutorBuilder( diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/SetExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/SetExecutor.java index 5394c8e51..c37d5f2ca 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/SetExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/SetExecutor.java @@ -31,7 +31,7 @@ public class SetExecutor extends RegularTaskExecutor { private final WorkflowFilter setFilter; - public static class SetExecutorBuilder extends RegularTaskExecutorBuilder { + public static class SetExecutorBuilder extends RegularTaskExecutorBuilder { private final WorkflowFilter setFilter; diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java index 7495c951a..a9a57269c 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java @@ -17,6 +17,7 @@ import io.serverlessworkflow.api.types.CatchErrors; import io.serverlessworkflow.api.types.ErrorFilter; +import io.serverlessworkflow.api.types.FlowDirective; import io.serverlessworkflow.api.types.Retry; import io.serverlessworkflow.api.types.RetryBackoff; import io.serverlessworkflow.api.types.RetryLimit; @@ -42,6 +43,7 @@ import io.serverlessworkflow.impl.executors.retry.RetryIntervalFunction; import java.time.Duration; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.concurrent.CompletableFuture; @@ -60,9 +62,10 @@ public class TryExecutor extends RegularTaskExecutor { private final Optional retryIntervalExecutor; private final Optional> attemptDuration; private final Optional> overallDuration; + private Optional catchTransition; private final String errorVariable; - public static class TryExecutorBuilder extends RegularTaskExecutorBuilder { + public static class TryExecutorBuilder extends RegularTaskExecutorBuilder { private final Optional whenFilter; private final Optional exceptFilter; @@ -72,6 +75,8 @@ public static class TryExecutorBuilder extends RegularTaskExecutorBuilder retryIntervalExecutor; private final Optional> attemptDuration; private final Optional> overallDuration; + private final FlowDirective catchDirective; + private Optional catchTransitionBuilder; private String errorVariable; protected TryExecutorBuilder( @@ -82,6 +87,7 @@ protected TryExecutorBuilder( this.errorFilter = buildErrorFilter(catchInfo.getErrors()); this.whenFilter = WorkflowUtils.optionalPredicate(application, catchInfo.getWhen()); this.exceptFilter = WorkflowUtils.optionalPredicate(application, catchInfo.getExceptWhen()); + this.catchDirective = catchInfo.getThen(); this.errorVariable = catchInfo.getAs(); List catchTaskDo = catchInfo.getDo(); this.catchTaskExecutor = @@ -99,6 +105,21 @@ protected TryExecutorBuilder( TaskExecutorHelper.createExecutorList(position, task.getTry(), definition, "try"); } + @Override + public void connect(Map> connections) { + super.connect(connections); + this.catchTransitionBuilder = + catchDirective != null + ? Optional.of(next(catchDirective, connections)) + : Optional.empty(); + } + + @Override + protected void buildTransition(TryExecutor instance) { + super.buildTransition(instance); + instance.catchTransition = catchTransitionBuilder.map(TransitionInfo::build); + } + private Optional resolveRetryPolicy(Retry retry) { RetryPolicy retryPolicy = null; if (retry != null) { @@ -253,6 +274,7 @@ private CompletableFuture handleException( .orElse(CompletableFuture.failedFuture(exception))) .thenCompose(model -> doIt(workflow, taskContext, model)); } + catchTransition.ifPresent(t -> taskContext.transition(t)); return completable; } else { return CompletableFuture.failedFuture(exception); @@ -319,4 +341,11 @@ && compareString(errorFilter.getTitle(), error.title()) private static boolean compareString(String one, String other) { return one == null || one.equals(other); } + + @Override + protected CompletableFuture execute( + WorkflowContext workflow, TaskContext taskContext) { + taskContext.transition(transition); + return internalExecute(workflow, taskContext).thenApply(node -> taskContext.rawOutput(node)); + } } 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 ba8b7a8c1..e6bc9366f 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 @@ -32,7 +32,8 @@ public class WaitExecutor extends RegularTaskExecutor { private final WorkflowValueResolver durationResolver; - public static class WaitExecutorBuilder extends RegularTaskExecutorBuilder { + public static class WaitExecutorBuilder + extends RegularTaskExecutorBuilder { private WorkflowValueResolver durationResolver; protected WaitExecutorBuilder( diff --git a/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java b/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java index 81c8866ce..5f8de1b8c 100644 --- a/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java +++ b/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java @@ -280,6 +280,7 @@ void testTimeout() throws IOException { @ValueSource( strings = { "workflows-samples/try-catch-match-when.yaml", + "workflows-samples/try-catch-match-then.yaml", "workflows-samples/try-catch-match-status.yaml", "workflows-samples/try-catch-match-details.yaml" }) diff --git a/impl/test/src/test/resources/workflows-samples/try-catch-match-then.yaml b/impl/test/src/test/resources/workflows-samples/try-catch-match-then.yaml new file mode 100644 index 000000000..551e94123 --- /dev/null +++ b/impl/test/src/test/resources/workflows-samples/try-catch-match-then.yaml @@ -0,0 +1,24 @@ +document: + dsl: '1.0.0' + namespace: test + name: try-catch-match-then + version: '0.1.0' +do: + - attemptTask: + try: + - failingTask: + raise: + error: + type: https://example.com/errors/transient + status: 503 + catch: + when: ${ .status == 503 } + then: endGracefully + - failWontHappen: + raise: + error: + type: https://serverlessworkflow.io/errors/runtime + status: 500 + - endGracefully: + set: + recovered: true From ec252934fb15d40848ebf8363479570e4d6ac013 Mon Sep 17 00:00:00 2001 From: Francisco Javier Tirado Sarti Date: Fri, 7 Aug 2026 19:14:56 +0200 Subject: [PATCH 2/2] [Fix #1601] Coverting retry and then Signed-off-by: Francisco Javier Tirado Sarti --- .../impl/executors/TryExecutor.java | 27 +++++++++----- .../impl/test/RetryTimeoutTest.java | 37 +++++++++++++++++++ .../try-catch-retry-inline-then.yaml | 37 +++++++++++++++++++ 3 files changed, 92 insertions(+), 9 deletions(-) create mode 100644 impl/test/src/test/resources/workflows-samples/try-catch-retry-inline-then.yaml diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java index a9a57269c..9d281a9d9 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/executors/TryExecutor.java @@ -263,18 +263,27 @@ private CompletableFuture handleException( TaskExecutorHelper.processTaskList( catchTaskExecutor.get(), workflow, Optional.of(taskContext), model)); } + if (retryIntervalExecutor.isPresent()) { completable = - completable - .thenCompose( - model -> - retryIntervalExecutor - .get() - .retry(workflow, taskContext, model) - .orElse(CompletableFuture.failedFuture(exception))) - .thenCompose(model -> doIt(workflow, taskContext, model)); + completable.thenCompose( + model -> { + Optional> retryCompletable = + retryIntervalExecutor.orElseThrow().retry(workflow, taskContext, model); + if (retryCompletable.isPresent()) { + return retryCompletable + .orElseThrow() + .thenCompose(innerModel -> doIt(workflow, taskContext, innerModel)); + } else if (catchTransition.isPresent()) { + taskContext.transition(catchTransition.orElseThrow()); + return CompletableFuture.completedFuture(model); + } else { + return CompletableFuture.failedFuture(exception); + } + }); + } else { + catchTransition.ifPresent(t -> taskContext.transition(t)); } - catchTransition.ifPresent(t -> taskContext.transition(t)); return completable; } else { return CompletableFuture.failedFuture(exception); diff --git a/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java b/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java index 5f8de1b8c..b59974d5b 100644 --- a/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java +++ b/impl/test/src/test/java/io/serverlessworkflow/impl/test/RetryTimeoutTest.java @@ -232,6 +232,43 @@ void testAttemptDurationRetry() throws IOException { assertThat(retryListener.taskRetried.get("do/0/tryGetPet/try/0/getPet")).isEqualTo((short) 1); } + @Test + void testAttemptRetryThen() throws IOException { + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue( + new MockResponse() + .setResponseCode(200) + .setHeader("Content-Type", "application/json") + .setBody(JsonUtils.mapper().writeValueAsString("{}"))); + CompletableFuture future = + app.workflowDefinition( + readWorkflowFromClasspath("workflows-samples/try-catch-retry-inline-then.yaml")) + .instance(Map.of()) + .start(); + assertThatThrownBy(() -> future.join()).hasCauseInstanceOf(WorkflowException.class); + } + + @Test + void testAttemptRetryThenCatch() throws IOException { + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + apiServer.enqueue(new MockResponse().setResponseCode(404)); + assertThat( + app.workflowDefinition( + readWorkflowFromClasspath("workflows-samples/try-catch-retry-inline-then.yaml")) + .instance(Map.of()) + .start() + .join() + .asMap() + .map(m -> m.get("recovered")) + .orElseThrow()) + .isEqualTo(true); + } + @Test void testAttemptDurationOverall() throws IOException { String result = "{\"name\":\"Luna\"}"; diff --git a/impl/test/src/test/resources/workflows-samples/try-catch-retry-inline-then.yaml b/impl/test/src/test/resources/workflows-samples/try-catch-retry-inline-then.yaml new file mode 100644 index 000000000..5f5df7d98 --- /dev/null +++ b/impl/test/src/test/resources/workflows-samples/try-catch-retry-inline-then.yaml @@ -0,0 +1,37 @@ +document: + dsl: '1.0.0' + namespace: test + name: try-catch-retry-inline-then + version: '0.1.0' +do: + - tryGetPet: + try: + - getPet: + call: http + with: + method: get + endpoint: http://localhost:9797 + redirect: true + catch: + errors: + with: + type: https://serverlessworkflow.io/spec/1.0.0/errors/communication + status: 404 + then: + endGracefully + retry: + delay: + milliseconds: 10 + backoff: + exponential: {} + limit: + attempt: + count: 5 + - failWontHappen: + raise: + error: + type: https://serverlessworkflow.io/errors/runtime + status: 500 + - endGracefully: + set: + recovered: true