From 1944cbede73095e85fd9055c97ac1b21743ab9c8 Mon Sep 17 00:00:00 2001 From: Francisco Javier Tirado Sarti Date: Mon, 27 Jul 2026 15:39:51 +0200 Subject: [PATCH] [Fix #1573] Send proper exception when there is validation error Also set the proper status when output validation or filtering fails. Signed-off-by: Francisco Javier Tirado Sarti --- .../impl/WorkflowApplication.java | 4 +- .../impl/WorkflowDefinition.java | 3 +- .../impl/WorkflowMutableInstance.java | 42 +++++++++------- .../impl/WorkflowUtils.java | 22 +++++++++ .../impl/executors/AbstractTaskExecutor.java | 10 ++-- .../impl/schema/AbstractSchemaValidator.java | 49 +++++++++++++++++++ .../impl/schema/SchemaValidator.java | 4 +- .../impl/test/HTTPWorkflowDefinitionTest.java | 16 +++--- .../jackson/schema/JsonSchemaValidator.java | 29 +++++------ .../io/serverlessworkflow/types/Errors.java | 1 + 10 files changed, 133 insertions(+), 47 deletions(-) create mode 100644 impl/core/src/main/java/io/serverlessworkflow/impl/schema/AbstractSchemaValidator.java 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 54dd4186d..116da6238 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowApplication.java @@ -219,7 +219,9 @@ private static final class EmptySchemaValidatorHolder { private final SchemaValidator NoValidation = new SchemaValidator() { @Override - public void validate(WorkflowModel node) {} + public Optional validate(WorkflowModel model) { + return Optional.empty(); + } }; @Override diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowDefinition.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowDefinition.java index a4065024d..1da55f6e2 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowDefinition.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowDefinition.java @@ -18,6 +18,7 @@ import static io.serverlessworkflow.impl.WorkflowUtils.buildWorkflowFilter; import static io.serverlessworkflow.impl.WorkflowUtils.getSchemaValidator; import static io.serverlessworkflow.impl.WorkflowUtils.safeClose; +import static io.serverlessworkflow.impl.WorkflowUtils.validationError; import io.serverlessworkflow.api.types.Input; import io.serverlessworkflow.api.types.Output; @@ -125,7 +126,7 @@ static WorkflowDefinition of(WorkflowApplication application, Workflow workflow, public WorkflowInstance instance(Object input) { WorkflowModel inputModel = application.modelFactory().fromAny(input); - inputSchemaValidator().ifPresent(v -> v.validate(inputModel)); + inputSchemaValidator().ifPresent(v -> validationError(v.validate(inputModel), definitionId)); return new WorkflowMutableInstance(this, application().idFactory().get(), inputModel); } 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 683518c87..c5a499e29 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java @@ -16,6 +16,7 @@ package io.serverlessworkflow.impl; import static io.serverlessworkflow.impl.LifecycleEventsUtils.publishEvent; +import static io.serverlessworkflow.impl.WorkflowUtils.validationError; import io.serverlessworkflow.impl.executors.TaskExecutorHelper; import io.serverlessworkflow.impl.lifecycle.WorkflowCancelledEvent; @@ -100,8 +101,9 @@ protected final CompletableFuture startExecution( .inputFilter() .map(f -> f.apply(workflowContext, null, input)) .orElse(input)) - .whenComplete(this::whenCompleted) - .thenApply(this::whenSuccess) + .whenComplete(this::setCompleteDate) + .thenApply(this::filterAndValidate) + .whenComplete(this::handleException) .thenCompose( model -> publishEvent( @@ -115,11 +117,8 @@ protected final CompletableFuture startExecution( return future; } - private void whenCompleted(WorkflowModel result, Throwable ex) { + private void setCompleteDate(WorkflowModel result, Throwable ex) { completedAt = Instant.now(); - if (ex != null) { - handleException(ex instanceof CompletionException ? ex = ex.getCause() : ex); - } } private void cleanUp(WorkflowModel result, Throwable ex) { @@ -131,23 +130,28 @@ private void cleanUp(WorkflowModel result, Throwable ex) { workflowContext.definition().removeInstance(this); } - private void handleException(Throwable ex) { - if (!(ex instanceof CancellationException)) { - status(WorkflowStatus.FAULTED); - publishEvent( - workflowContext, l -> l.onWorkflowFailed(new WorkflowFailedEvent(workflowContext, ex))); + 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 WorkflowModel whenSuccess(WorkflowModel node) { + private WorkflowModel filterAndValidate(WorkflowModel model) { + WorkflowDefinition definition = workflowContext.definition(); WorkflowModel output = - workflowContext - .definition() - .outputFilter() - .map(f -> f.apply(workflowContext, null, node)) - .orElse(node); - workflowContext.definition().outputSchemaValidator().ifPresent(v -> v.validate(output)); - status(WorkflowStatus.COMPLETED); + definition.outputFilter().map(f -> f.apply(workflowContext, null, model)).orElse(model); + definition + .outputSchemaValidator() + .ifPresent(v -> validationError(v.validate(output), workflowContext)); return output; } diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java index 5ff4ea0d6..f1e3e35f4 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowUtils.java @@ -318,4 +318,26 @@ public static Optional loadFirst(Class service .sorted() .findFirst(); } + + public static void validationError( + Optional errorBuilder, TaskContextData taskContext) { + validationError(errorBuilder, taskContext.position().jsonPointer().toString()); + } + + public static void validationError( + Optional errorBuilder, WorkflowDefinitionId definition) { + validationError(errorBuilder, definition.toString()); + } + + public static void validationError( + Optional errorBuilder, WorkflowContextData workflowContext) { + validationError(errorBuilder, workflowContext.instanceData().id()); + } + + private static void validationError( + Optional errorBuilder, String instance) { + if (errorBuilder.isPresent()) { + throw new WorkflowException(errorBuilder.orElseThrow().instance(instance).build()); + } + } } 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 591b3f503..af335f7a9 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 @@ -18,6 +18,7 @@ import static io.serverlessworkflow.impl.LifecycleEventsUtils.publishEvent; import static io.serverlessworkflow.impl.WorkflowUtils.buildWorkflowFilter; import static io.serverlessworkflow.impl.WorkflowUtils.getSchemaValidator; +import static io.serverlessworkflow.impl.WorkflowUtils.validationError; import io.serverlessworkflow.api.types.Export; import io.serverlessworkflow.api.types.FlowDirective; @@ -232,7 +233,8 @@ public CompletableFuture apply( }) .thenApply( t -> { - inputSchemaValidator.ifPresent(s -> s.validate(t.rawInput())); + inputSchemaValidator.ifPresent( + s -> validationError(s.validate(t.rawInput()), t)); inputProcessor.ifPresent( p -> taskContext.input(p.apply(workflowContext, t, t.rawInput()))); return t; @@ -252,12 +254,14 @@ public CompletableFuture apply( t -> { outputProcessor.ifPresent( p -> t.output(p.apply(workflowContext, t, t.rawOutput()))); - outputSchemaValidator.ifPresent(s -> s.validate(t.output())); + outputSchemaValidator.ifPresent( + s -> validationError(s.validate(t.output()), t)); contextProcessor.ifPresent( p -> workflowContext.context( p.apply(workflowContext, t, workflowContext.context()))); - contextSchemaValidator.ifPresent(s -> s.validate(workflowContext.context())); + contextSchemaValidator.ifPresent( + s -> validationError(s.validate(workflowContext.context()), t)); t.completedAt(Instant.now()); return t; }) diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/schema/AbstractSchemaValidator.java b/impl/core/src/main/java/io/serverlessworkflow/impl/schema/AbstractSchemaValidator.java new file mode 100644 index 000000000..be688443b --- /dev/null +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/schema/AbstractSchemaValidator.java @@ -0,0 +1,49 @@ +/* + * 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.schema; + +import io.serverlessworkflow.impl.WorkflowError; +import io.serverlessworkflow.impl.WorkflowModel; +import io.serverlessworkflow.types.Errors; +import java.util.Optional; + +public abstract class AbstractSchemaValidator implements SchemaValidator { + + private final Class validationClass; + + protected AbstractSchemaValidator(Class validationClass) { + this.validationClass = validationClass; + } + + @Override + public Optional validate(WorkflowModel model) { + return validate( + model + .as(validationClass) + .orElseThrow( + () -> + new IllegalStateException( + "Model cannot be converted to proper schema validation class " + + validationClass))) + .map( + s -> + WorkflowError.error(Errors.VALIDATION.toString(), Errors.VALIDATION.status()) + .title("Schema validation errors") + .details(s)); + } + + protected abstract Optional validate(T object); +} diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/schema/SchemaValidator.java b/impl/core/src/main/java/io/serverlessworkflow/impl/schema/SchemaValidator.java index fa66676b9..fa52364ec 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/schema/SchemaValidator.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/schema/SchemaValidator.java @@ -15,8 +15,10 @@ */ package io.serverlessworkflow.impl.schema; +import io.serverlessworkflow.impl.WorkflowError; import io.serverlessworkflow.impl.WorkflowModel; +import java.util.Optional; public interface SchemaValidator { - void validate(WorkflowModel node); + Optional validate(WorkflowModel model); } diff --git a/impl/test/src/test/java/io/serverlessworkflow/impl/test/HTTPWorkflowDefinitionTest.java b/impl/test/src/test/java/io/serverlessworkflow/impl/test/HTTPWorkflowDefinitionTest.java index e94646b76..6521074cd 100644 --- a/impl/test/src/test/java/io/serverlessworkflow/impl/test/HTTPWorkflowDefinitionTest.java +++ b/impl/test/src/test/java/io/serverlessworkflow/impl/test/HTTPWorkflowDefinitionTest.java @@ -21,7 +21,10 @@ import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import io.serverlessworkflow.impl.WorkflowApplication; +import io.serverlessworkflow.impl.WorkflowError; +import io.serverlessworkflow.impl.WorkflowException; import io.serverlessworkflow.impl.WorkflowModel; +import io.serverlessworkflow.types.Errors; import java.io.IOException; import java.nio.charset.StandardCharsets; import java.util.Map; @@ -373,18 +376,19 @@ void callHttpPut_should_contain_firstName_with_john() throws Exception { } @Test - void testWrongSchema_should_throw_illegal_argument() { - IllegalArgumentException exception = + void testWrongSchema_should_throw_exception() { + WorkflowException exception = catchThrowableOfType( - IllegalArgumentException.class, + WorkflowException.class, () -> appl.workflowDefinition( readWorkflowFromClasspath( "workflows-samples/call-http-query-parameters.yaml")) .instance(Map.of())); - assertThat(exception) - .isNotNull() - .hasMessageContaining("There are JsonSchema validation errors"); + assertThat(exception).isNotNull(); + WorkflowError error = exception.getWorkflowError(); + assertThat(error.status()).isEqualTo(400); + assertThat(error.type()).isEqualTo(Errors.VALIDATION.toString()); } @Test diff --git a/impl/validation/src/main/java/io/serverlessworkflow/impl/jackson/schema/JsonSchemaValidator.java b/impl/validation/src/main/java/io/serverlessworkflow/impl/jackson/schema/JsonSchemaValidator.java index d01047d08..be8f6053a 100644 --- a/impl/validation/src/main/java/io/serverlessworkflow/impl/jackson/schema/JsonSchemaValidator.java +++ b/impl/validation/src/main/java/io/serverlessworkflow/impl/jackson/schema/JsonSchemaValidator.java @@ -20,32 +20,29 @@ import com.networknt.schema.Schema; import com.networknt.schema.SchemaRegistry; import com.networknt.schema.SpecificationVersion; -import io.serverlessworkflow.impl.WorkflowModel; -import io.serverlessworkflow.impl.schema.SchemaValidator; +import io.serverlessworkflow.impl.schema.AbstractSchemaValidator; import java.util.Collection; +import java.util.Optional; +import java.util.stream.Collectors; -public class JsonSchemaValidator implements SchemaValidator { +public class JsonSchemaValidator extends AbstractSchemaValidator { private final Schema schemaObject; public JsonSchemaValidator(JsonNode jsonNode) { + super(JsonNode.class); this.schemaObject = SchemaRegistry.withDefaultDialect(SpecificationVersion.DRAFT_7).getSchema(jsonNode); } @Override - public void validate(WorkflowModel node) { - Collection report = - schemaObject.validate( - node.as(JsonNode.class) - .orElseThrow( - () -> - new IllegalArgumentException( - "Default schema validator requires WorkflowModel to support conversion to json node"))); - if (!report.isEmpty()) { - StringBuilder sb = new StringBuilder("There are JsonSchema validation errors:"); - report.forEach(m -> sb.append(System.lineSeparator()).append(m.getMessage())); - throw new IllegalArgumentException(sb.toString()); - } + protected Optional validate(JsonNode node) { + Collection report = schemaObject.validate(node); + return report.isEmpty() + ? Optional.empty() + : Optional.of( + report.stream() + .map(Error::getMessage) + .collect(Collectors.joining(System.lineSeparator()))); } } diff --git a/types/src/main/java/io/serverlessworkflow/types/Errors.java b/types/src/main/java/io/serverlessworkflow/types/Errors.java index e7b090569..2c2f25cb3 100644 --- a/types/src/main/java/io/serverlessworkflow/types/Errors.java +++ b/types/src/main/java/io/serverlessworkflow/types/Errors.java @@ -93,4 +93,5 @@ public String toString() { public static final Standard DATA = new Standard("data", 422); public static final Standard TIMEOUT = new Standard("timeout", 408); public static final Standard NOT_IMPLEMENTED = new Standard("not-implemented", 501); + public static final Standard VALIDATION = new Standard("validation", 400); }