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 @@ -219,7 +219,9 @@ private static final class EmptySchemaValidatorHolder {
private final SchemaValidator NoValidation =
new SchemaValidator() {
@Override
public void validate(WorkflowModel node) {}
public Optional<WorkflowError.Builder> validate(WorkflowModel model) {
return Optional.empty();
}
};

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -100,8 +101,9 @@ protected final CompletableFuture<WorkflowModel> 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(
Expand All @@ -115,11 +117,8 @@ protected final CompletableFuture<WorkflowModel> 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) {
Expand All @@ -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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -318,4 +318,26 @@ public static <T extends ServicePriority> Optional<T> loadFirst(Class<T> service
.sorted()
.findFirst();
}

public static void validationError(
Optional<WorkflowError.Builder> errorBuilder, TaskContextData taskContext) {
validationError(errorBuilder, taskContext.position().jsonPointer().toString());
}

public static void validationError(
Optional<WorkflowError.Builder> errorBuilder, WorkflowDefinitionId definition) {
validationError(errorBuilder, definition.toString());
}

public static void validationError(
Optional<WorkflowError.Builder> errorBuilder, WorkflowContextData workflowContext) {
validationError(errorBuilder, workflowContext.instanceData().id());
}

private static void validationError(
Optional<WorkflowError.Builder> errorBuilder, String instance) {
if (errorBuilder.isPresent()) {
throw new WorkflowException(errorBuilder.orElseThrow().instance(instance).build());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -232,7 +233,8 @@ public CompletableFuture<TaskContext> 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;
Comment thread
fjtirado marked this conversation as resolved.
Expand All @@ -252,12 +254,14 @@ public CompletableFuture<TaskContext> 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;
})
Expand Down
Original file line number Diff line number Diff line change
@@ -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<T> implements SchemaValidator {

private final Class<T> validationClass;

protected AbstractSchemaValidator(Class<T> validationClass) {
this.validationClass = validationClass;
}

@Override
public Optional<WorkflowError.Builder> 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<String> validate(T object);
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<WorkflowError.Builder> validate(WorkflowModel model);
}
Comment thread
fjtirado marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<JsonNode> {

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<Error> 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<String> validate(JsonNode node) {
Collection<Error> report = schemaObject.validate(node);
return report.isEmpty()
? Optional.empty()
: Optional.of(
report.stream()
.map(Error::getMessage)
.collect(Collectors.joining(System.lineSeparator())));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Loading