From 985d2fb67f6007ef761fc6eb91674e9202f3b479 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Wed, 30 Jul 2025 16:59:17 +0200 Subject: [PATCH] HumanInteraction node added --- .../manager/WorkflowManagerApplication.java | 18 +- .../manager/controllers/FlowsController.java | 7 +- .../{nodes => }/NodeDefinitionController.java | 46 +++-- .../cnr/isti/workflow/manager/dto/Flow.java | 44 +++-- .../cnr/isti/workflow/manager/dto/Node.java | 5 +- .../workflow/manager/entities/FlowEntity.java | 8 +- .../HumanInteractionNodeDefinition.java | 40 ++++ .../definitions/InputNodeDefinition.java | 21 ++- .../nodes/definitions/NodeDefinition.java | 12 +- .../definitions/OutputNodeDefinition.java | 19 +- .../definitions/SystemNodeDefinition.java | 25 +++ .../nodes/definitions/UserNodeDefinition.java | 9 +- .../manager/executors/ExecutionObject.java | 3 - .../manager/executors/NoOpExecutor.java | 2 - .../executors/ai/models/GeminiModel.java | 3 +- .../manager/model/ExecutionContext.java | 24 ++- .../manager/model/ExecutorDescriptor.java | 1 - .../manager/model/InputReadyListener.java | 4 +- .../manager/model/steps/ExecutionStep.java | 159 +--------------- .../manager/model/steps/InputStep.java | 4 +- .../manager/model/steps/InteractiveStep.java | 78 ++++++++ .../manager/model/steps/OutputStep.java | 7 + .../workflow/manager/model/steps/Step.java | 6 - .../manager/model/steps/UserDefinedStep.java | 171 ++++++++++++++++++ .../model/types/ParameterDefinition.java | 1 - .../InputDefinitionRepository.java | 10 - .../NodeDefinitionRepository.java | 11 +- .../OutputDefinitionRepository.java | 9 - .../SystemNodeDefinitionRepository.java | 15 ++ .../UserNodeDefinitionRepository.java | 10 +- .../manager/services/ExecutionService.java | 62 ++++--- .../manager/services/ImportComponent.java | 3 +- .../manager/services/TransformerService.java | 54 ++++-- .../isti/workflow/manager/Configurations.java | 1 - .../manager/ExecutionServiceTest.java | 35 ++++ .../cnr/isti/workflow/manager/FlowTest.java | 11 +- .../workflow/manager/TestFlowCreator.java | 6 +- .../manager/repositories/NodeDefRepoTest.java | 14 +- src/test/resources/interactionFlow.json | 168 +++++++++++++++++ 39 files changed, 797 insertions(+), 329 deletions(-) rename src/main/java/it/cnr/isti/workflow/manager/controllers/{nodes => }/NodeDefinitionController.java (51%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/HumanInteractionNodeDefinition.java rename src/main/java/it/cnr/isti/workflow/manager/{model/types => entities}/nodes/definitions/InputNodeDefinition.java (64%) rename src/main/java/it/cnr/isti/workflow/manager/{model/types => entities}/nodes/definitions/NodeDefinition.java (83%) rename src/main/java/it/cnr/isti/workflow/manager/{model/types => entities}/nodes/definitions/OutputNodeDefinition.java (64%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/SystemNodeDefinition.java rename src/main/java/it/cnr/isti/workflow/manager/{model/types => entities}/nodes/definitions/UserNodeDefinition.java (87%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/model/steps/InteractiveStep.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/model/steps/UserDefinedStep.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/repositories/InputDefinitionRepository.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/repositories/OutputDefinitionRepository.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/repositories/SystemNodeDefinitionRepository.java create mode 100644 src/test/resources/interactionFlow.json diff --git a/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java b/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java index 2dbcd20..c81fe50 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java +++ b/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java @@ -7,13 +7,13 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Import; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.OutputNodeDefinition; import it.cnr.isti.workflow.manager.executors.ai.models.AIModel; import it.cnr.isti.workflow.manager.model.types.IOType; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; import it.cnr.isti.workflow.manager.services.ImportComponent; @@ -36,14 +36,20 @@ public class WorkflowManagerApplication { return args -> { NodeDefinition textInput = InputNodeDefinition.builder().name("Text Input") - .createdBy("admin").category("Inputs").outputType(IOType.TEXT) + .outputType(IOType.TEXT) .build(); NodeDefinition textOutput = OutputNodeDefinition.builder().name("Text Output") - .createdBy("admin").category("Outputs").inputType(IOType.TEXT) + .inputType(IOType.TEXT) + .build(); + + HumanInteractionNodeDefinition humanInteractionNodeDefinition = HumanInteractionNodeDefinition.builder() + .name("Human Interaction") + .simulable(false) .build(); repository.save(textInput); repository.save(textOutput); + repository.save(humanInteractionNodeDefinition); importComponent.start(); }; diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowsController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowsController.java index 7488ca7..b006d81 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowsController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowsController.java @@ -90,9 +90,8 @@ public class FlowsController { if (flow.getId() != null) throw new WebException(HttpStatus.BAD_REQUEST, "on creation flow id must be null"); flow.setCreatedBy(userDetails.getUsername()); - Flow savedFlow = flowRepository.save((FlowEntity)flow); - - logger.info("flow created {}", flow.getId()); + FlowEntity savedFlow = flowRepository.save(FlowEntity.from(flow)); + logger.info("flow created {}", savedFlow.getId()); return savedFlow; } @@ -113,7 +112,7 @@ public class FlowsController { //validateFlow(flow); if (!id.equals(flow.getId())) throw new WebException(HttpStatus.BAD_REQUEST, "Flow id does not match"); - Flow savedFlow = flowRepository.findById(id).orElseThrow(() -> new WebException(HttpStatus.NOT_FOUND,"Flow not found")); + FlowEntity savedFlow = flowRepository.findById(id).orElseThrow(() -> new WebException(HttpStatus.NOT_FOUND,"Flow not found")); if (!savedFlow.getCreatedBy().equals(userDetails.getUsername())) throw new WebException(HttpStatus.FORBIDDEN, "Flow not created by user"); diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/nodes/NodeDefinitionController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java similarity index 51% rename from src/main/java/it/cnr/isti/workflow/manager/controllers/nodes/NodeDefinitionController.java rename to src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java index 1b18b74..0733da6 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/nodes/NodeDefinitionController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java @@ -1,4 +1,4 @@ -package it.cnr.isti.workflow.manager.controllers.nodes; +package it.cnr.isti.workflow.manager.controllers; import java.util.List; @@ -11,11 +11,13 @@ import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import io.swagger.v3.oas.annotations.security.SecurityRequirement; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.executors.Executor; import it.cnr.isti.workflow.manager.model.ExecutorDescriptor; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; +import it.cnr.isti.workflow.manager.repositories.SystemNodeDefinitionRepository; +import it.cnr.isti.workflow.manager.repositories.UserNodeDefinitionRepository; @RestController @RequestMapping("/types") @@ -26,26 +28,44 @@ public class NodeDefinitionController { @Autowired private NodeDefinitionRepository repository; + @Autowired + private UserNodeDefinitionRepository userNodeDefinitionRepository; + + @Autowired + private SystemNodeDefinitionRepository systemNodeDefinitionRepository; + @Autowired private List executors; //@SecurityRequirement(name = "bearerAuth") - @GetMapping("/nodes") - public List getNodes() { - return (List) repository.findAll(); + @GetMapping("/nodes/system") + public List getSystemNodes() { + return (List) systemNodeDefinitionRepository.findAll(); } - - @GetMapping("/nodes/types") - public List getStoredNodeTypes() { - return repository.getUsedTypes(); + @GetMapping("/nodes/user") + public List getUserNodes() { + return (List) userNodeDefinitionRepository.findAll(); + + } + + @GetMapping("/nodes/system/categories") + public List getSystemNodeCategories() { + return systemNodeDefinitionRepository.getCategories(); + + } + + + @GetMapping("/nodes/user/types") + public List getUserNodeTypes() { + return userNodeDefinitionRepository.getUsedTypes(); } - @GetMapping("/nodes/categories") - public List getStoredNodeCategories() { - return repository.getCategories(); + @GetMapping("/nodes/user/categories") + public List getUserNodeCategories() { + return userNodeDefinitionRepository.getUsedCategories(); } @SecurityRequirement(name = "bearerAuth") diff --git a/src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java b/src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java index c57e5ea..1719468 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java @@ -2,39 +2,61 @@ package it.cnr.isti.workflow.manager.dto; import java.util.List; import com.fasterxml.jackson.annotation.JsonAlias; -import lombok.AllArgsConstructor; +import com.fasterxml.jackson.databind.annotation.JsonDeserialize; + +import lombok.AccessLevel; import lombok.Builder; -import lombok.Builder.Default; import lombok.Data; import lombok.NoArgsConstructor; import lombok.NonNull; import lombok.ToString; @Data -@AllArgsConstructor -@NoArgsConstructor -@Builder +@NoArgsConstructor(access = AccessLevel.PROTECTED) @ToString public class Flow { + @Builder + protected Flow(String id, String createdBy, boolean isPublic, String name, String description, List nodes, + List connections) { + this.id = id; + this.createdBy = createdBy; + this.isPublic = isPublic; + this.name = name; + this.description = description; + this.nodes = nodes; + this.connections = connections; + if (nodes != null) { + this.loadNodeDefinitions(); + } + } + String id; @NonNull String createdBy; - @Default @JsonAlias("public") - boolean isPublic =false; - + boolean isPublic = false; + @NonNull String name; String description; - @Builder.Default - List nodes = List.of(); + protected List nodes = List.of(); + + public void setNodes(List nodes) { + this.nodes = nodes; + if (nodes != null) { + this.loadNodeDefinitions(); + } + + } - @Builder.Default List connections = List.of(); + protected void loadNodeDefinitions() { + this.getNodes().forEach(Node::resolveNodeDefinition); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/dto/Node.java b/src/main/java/it/cnr/isti/workflow/manager/dto/Node.java index e69aab0..95a354d 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/dto/Node.java +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Node.java @@ -3,10 +3,9 @@ package it.cnr.isti.workflow.manager.dto; import java.util.List; import java.util.Map; -import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; import it.cnr.isti.workflow.manager.services.RepositoryHolder; import lombok.AllArgsConstructor; import lombok.Builder; @@ -46,7 +45,7 @@ public class Node { String type; public void resolveNodeDefinition() { - if (type != null) { + if (type != null && nodeDefinition == null) { this.nodeDefinition = RepositoryHolder.getNodeDefinition() .findById(type) .orElse(null); diff --git a/src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java b/src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java index bbe5bd2..22dbb6a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java @@ -14,6 +14,8 @@ import jakarta.persistence.Convert; import jakarta.persistence.Entity; import jakarta.persistence.Id; import jakarta.persistence.PostLoad; +import jakarta.persistence.PostPersist; +import jakarta.persistence.PostUpdate; import lombok.EqualsAndHashCode; @Entity @@ -106,7 +108,7 @@ public class FlowEntity extends Flow { @Override public void setNodes(List nodes) { - super.setNodes(nodes); + this.nodes = nodes; } @Convert(converter = ConnectionListConverter.class) @@ -121,8 +123,10 @@ public class FlowEntity extends Flow { super.setConnections(connections); } + @PostPersist + @PostUpdate @PostLoad public void postLoad() { - this.getNodes().forEach(Node::resolveNodeDefinition); + super.loadNodeDefinitions(); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/HumanInteractionNodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/HumanInteractionNodeDefinition.java new file mode 100644 index 0000000..ac608bf --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/HumanInteractionNodeDefinition.java @@ -0,0 +1,40 @@ +package it.cnr.isti.workflow.manager.entities.nodes.definitions; + +import java.util.Map; + +import com.fasterxml.jackson.annotation.JsonTypeName; + +import it.cnr.isti.workflow.manager.model.types.IOType; +import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; +import it.cnr.isti.workflow.manager.model.types.ParameterType; +import jakarta.persistence.Entity; +import lombok.Builder; +import lombok.NoArgsConstructor; + +@Entity +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@JsonTypeName("system") +public class HumanInteractionNodeDefinition extends SystemNodeDefinition { + + public static final String PARAMETER_NAME = "action_description"; + public static final String INPUT_NAME = "input"; + public static final String OUTPUT_NAME = "output"; + + @Builder + public HumanInteractionNodeDefinition(String name, boolean simulable) { + super(name, simulable); + this.setCategory("humanInteraction"); + this.outputs = Map.of(OUTPUT_NAME, IOType.TEXT); + this.inputs = Map.of(INPUT_NAME, IOType.TEXT); + + this.setColor("black"); + this.runtimeParameters = Map.of(PARAMETER_NAME, ParameterDefinition.builder() + .name(PARAMETER_NAME) + .label("Human Readable Action Description") + .description("Description of the user action to perform") + .required(true) + .type(ParameterType.Text) + .build()); + } + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/InputNodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/InputNodeDefinition.java similarity index 64% rename from src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/InputNodeDefinition.java rename to src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/InputNodeDefinition.java index 16edfe3..aa8b6f9 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/InputNodeDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/InputNodeDefinition.java @@ -1,37 +1,40 @@ -package it.cnr.isti.workflow.manager.model.types.nodes.definitions; +package it.cnr.isti.workflow.manager.entities.nodes.definitions; import java.util.Map; +import com.fasterxml.jackson.annotation.JsonTypeName; + import it.cnr.isti.workflow.manager.model.types.IOType; import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; import it.cnr.isti.workflow.manager.model.types.ParameterType; -import jakarta.persistence.DiscriminatorValue; import jakarta.persistence.Entity; import lombok.Builder; import lombok.NoArgsConstructor; @Entity @NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) -@DiscriminatorValue("2") -public class InputNodeDefinition extends NodeDefinition { +@JsonTypeName("system") +public class InputNodeDefinition extends SystemNodeDefinition { public static final String INPUT_PARAMETER_KEY ="name"; public static final String OUTPUT_NAME = "value"; + public static final String DEFAULT_VALUE = "output"; + @Builder - public InputNodeDefinition(String name, String createdBy, String category, IOType outputType) { + public InputNodeDefinition(String name, IOType outputType, boolean simulable) { super(); + this.simulable = simulable; this.setName(name); - this.setCreatedBy(createdBy); - this.setCategory(category); + this.setCategory("input"); this.outputs = Map.of(OUTPUT_NAME, outputType); - this.setColor("blue"); + this.setColor("green"); this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder() .name(INPUT_PARAMETER_KEY) .label("Input Name") .description("Input name") - .defaultValue("output") + .defaultValue(DEFAULT_VALUE) .type(ParameterType.String) .build()); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/NodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/NodeDefinition.java similarity index 83% rename from src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/NodeDefinition.java rename to src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/NodeDefinition.java index 1fdb97a..83b6c7e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/NodeDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/NodeDefinition.java @@ -1,8 +1,10 @@ -package it.cnr.isti.workflow.manager.model.types.nodes.definitions; +package it.cnr.isti.workflow.manager.entities.nodes.definitions; import java.util.HashMap; import java.util.Map; +import org.springframework.boot.autoconfigure.security.SecurityProperties.User; + import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; @@ -35,13 +37,7 @@ import lombok.experimental.SuperBuilder; @Inheritance(strategy = InheritanceType.SINGLE_TABLE) @DiscriminatorColumn(name="node_type", length = 1, discriminatorType = DiscriminatorType.INTEGER) - @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type", defaultImpl = UserNodeDefinition.class) -@JsonSubTypes({ - @JsonSubTypes.Type(value = InputNodeDefinition.class, name = "input"), - @JsonSubTypes.Type(value = UserNodeDefinition.class, name = "user"), - @JsonSubTypes.Type(value = OutputNodeDefinition.class, name = "output"), - -}) +@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type", defaultImpl = UserNodeDefinition.class) public abstract class NodeDefinition { @Id diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/OutputNodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/OutputNodeDefinition.java similarity index 64% rename from src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/OutputNodeDefinition.java rename to src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/OutputNodeDefinition.java index ea26521..771af5f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/OutputNodeDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/OutputNodeDefinition.java @@ -1,34 +1,33 @@ -package it.cnr.isti.workflow.manager.model.types.nodes.definitions; +package it.cnr.isti.workflow.manager.entities.nodes.definitions; import java.util.Map; -import it.cnr.isti.workflow.manager.model.Validations; +import com.fasterxml.jackson.annotation.JsonTypeName; + import it.cnr.isti.workflow.manager.model.types.IOType; import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; import it.cnr.isti.workflow.manager.model.types.ParameterType; -import it.cnr.isti.workflow.manager.model.types.Validation; -import jakarta.persistence.DiscriminatorValue; import jakarta.persistence.Entity; import lombok.Builder; import lombok.NoArgsConstructor; @Entity @NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) -@DiscriminatorValue("3") -public class OutputNodeDefinition extends NodeDefinition { +@JsonTypeName("system") +public class OutputNodeDefinition extends SystemNodeDefinition { public static final String INPUT_PARAMETER_KEY ="name"; public static final String INPUT_NAME = "value"; @Builder - public OutputNodeDefinition(String name, String createdBy, String category, IOType inputType) { + public OutputNodeDefinition(String name, IOType inputType, boolean simulable) { super(); + this.simulable = simulable; this.setName(name); - this.setCreatedBy(createdBy); - this.setCategory(category); + this.setCategory("output"); this.inputs = Map.of(INPUT_NAME, inputType); - this.setColor("blue"); + this.setColor("red"); this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder() .name(INPUT_PARAMETER_KEY) .label("Output Name") diff --git a/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/SystemNodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/SystemNodeDefinition.java new file mode 100644 index 0000000..796e0b0 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/SystemNodeDefinition.java @@ -0,0 +1,25 @@ +package it.cnr.isti.workflow.manager.entities.nodes.definitions; + + + +import com.fasterxml.jackson.annotation.JsonProperty; + +import jakarta.persistence.DiscriminatorValue; +import jakarta.persistence.Entity; +import lombok.NoArgsConstructor; + +@Entity +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@DiscriminatorValue("1") +public abstract class SystemNodeDefinition extends NodeDefinition { + + @JsonProperty + boolean simulable = false; + + public SystemNodeDefinition(String name, boolean simulable) { + super(); + this.simulable = simulable; + this.setName(name); + this.setCreatedBy("system"); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/UserNodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/UserNodeDefinition.java similarity index 87% rename from src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/UserNodeDefinition.java rename to src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/UserNodeDefinition.java index 29a6c7e..828646a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/nodes/definitions/UserNodeDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/nodes/definitions/UserNodeDefinition.java @@ -1,6 +1,10 @@ -package it.cnr.isti.workflow.manager.model.types.nodes.definitions; +package it.cnr.isti.workflow.manager.entities.nodes.definitions; import java.util.Map; + +import com.fasterxml.jackson.annotation.JsonTypeInfo; +import com.fasterxml.jackson.annotation.JsonTypeName; + import it.cnr.isti.workflow.manager.model.types.Translator; import it.cnr.isti.workflow.manager.repositories.converters.InputTranslatorConverter; import it.cnr.isti.workflow.manager.repositories.converters.MapsConverter; @@ -23,7 +27,8 @@ import lombok.experimental.SuperBuilder; @SuperBuilder @Data @EqualsAndHashCode(callSuper = true) -@DiscriminatorValue("1") +@DiscriminatorValue("2") +@JsonTypeName("user") public class UserNodeDefinition extends NodeDefinition { @NonNull diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java index d166623..d33ec82 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ExecutionObject.java @@ -5,9 +5,6 @@ import java.util.List; import java.util.Map; import java.util.UUID; -import com.fasterxml.jackson.annotation.JsonIgnore; - -import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.model.ExecutionContext; import it.cnr.isti.workflow.manager.model.FieldKey; import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java index d037d75..87a7e20 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java @@ -1,8 +1,6 @@ package it.cnr.isti.workflow.manager.executors; import org.springframework.stereotype.Service; - -import java.lang.reflect.Parameter; import java.util.List; import java.util.Map; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/models/GeminiModel.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/models/GeminiModel.java index c3d735e..3cd4929 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/models/GeminiModel.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/models/GeminiModel.java @@ -1,7 +1,6 @@ package it.cnr.isti.workflow.manager.executors.ai.models; import java.util.Map; -import org.slf4j.Logger; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; @@ -11,7 +10,7 @@ import it.cnr.isti.workflow.manager.executors.ai.services.GeminiService; @Service("google-gemini") public class GeminiModel implements AIModel { - private static final Logger log = org.slf4j.LoggerFactory.getLogger(GeminiModel.class); + //private static final Logger log = org.slf4j.LoggerFactory.getLogger(GeminiModel.class); private GeminiService geminiService; diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java index 8fd89e4..50b9932 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java @@ -13,20 +13,23 @@ import lombok.NoArgsConstructor; public class ExecutionContext { public enum Status { - CREATED(true, false), - INITIALIZING(true, false), - READY(true, false), - RUNNING(false, false), - SUCCESS(false, true), - ERROR(false, true); + CREATED(true, false, false), + INITIALIZING(true, false, false), + READY(true, false, false), + RUNNING(false, false, false), + WAITING(false, false, true), + SUCCESS(false, true, false), + ERROR(false, true, false); private boolean initState; private boolean finalState; + private boolean waitingState; - Status(boolean initState, boolean finalState) { + Status(boolean initState, boolean finalState, boolean waitingState) { this.initState = initState; this.finalState = finalState; + this.waitingState = waitingState; } public boolean isInitState() { @@ -37,8 +40,12 @@ public class ExecutionContext { return finalState; } + public boolean isWaitingState() { + return waitingState ; + } + public boolean isRunningState() { - return !finalState && !initState; + return !finalState && !initState && !waitingState; } } @@ -51,6 +58,7 @@ public class ExecutionContext { Long endTime = null; List stepsUnderExecution = new ArrayList<>(); + List waitingSteps = new ArrayList<>(); Map errors = new HashMap<>(); List warnings = new ArrayList<>(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java index 305e123..f64a59c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java @@ -2,7 +2,6 @@ package it.cnr.isti.workflow.manager.model; import java.util.List; -import it.cnr.isti.workflow.manager.executors.Executor; import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; import lombok.AllArgsConstructor; import lombok.Builder; diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/InputReadyListener.java b/src/main/java/it/cnr/isti/workflow/manager/model/InputReadyListener.java index 9e16038..7a4a68b 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/InputReadyListener.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/InputReadyListener.java @@ -1,10 +1,8 @@ package it.cnr.isti.workflow.manager.model; -import java.util.Map; - public interface InputReadyListener { void inputReady(String key, Object value); - Map getParentOutputsToInputMapping(); + void addParentOutputsToInputMapping(String parentOutputName, String inputName); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/ExecutionStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/ExecutionStep.java index c198029..dd528db 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/steps/ExecutionStep.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/ExecutionStep.java @@ -1,165 +1,12 @@ package it.cnr.isti.workflow.manager.model.steps; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.stream.Collectors; -import org.slf4j.Logger; import it.cnr.isti.workflow.manager.dto.Point; -import com.fasterxml.jackson.annotation.JsonIgnore; -import it.cnr.isti.workflow.manager.exceptions.ExecutionException; -import it.cnr.isti.workflow.manager.executors.Executor; -import it.cnr.isti.workflow.manager.model.ExecutionContext; -import it.cnr.isti.workflow.manager.model.InputReadyListener; -import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; -import it.cnr.isti.workflow.manager.model.types.IOType; -import it.cnr.isti.workflow.manager.model.types.Translator; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; -import lombok.Builder; -import lombok.Data; -import lombok.EqualsAndHashCode; -import lombok.NonNull; -import lombok.ToString; -@Data -@ToString -@EqualsAndHashCode(callSuper = true) -public class ExecutionStep extends Step implements InputReadyListener, ExecutableStep { +public abstract class ExecutionStep extends Step { - private static final Logger log = org.slf4j.LoggerFactory.getLogger(ExecutionStep.class); - - @NonNull - String name; - - Map runtimeParameters; - Map fixedParameters; - - @JsonIgnore - @NonNull - ExecutionContext context; - - Map inputTypes; - Map outputTypes; - - @JsonIgnore - Executor executor; - - @JsonIgnore - Map parentOutputsToInputMapping = new HashMap<>(); - - @JsonIgnore - Map inputs = new HashMap<>(); - - @JsonIgnore - @NonNull - UserNodeDefinition nodeDefinition; - - @Builder - public ExecutionStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map runtimeParameters, - Map fixedParameters, - UserNodeDefinition nodeDefinition, - Executor executor, Point position, @NonNull Map inputTypes, - @NonNull Map outputTypes) { + public ExecutionStep(String id, Point position) { super(id, position); - this.name = name; - this.runtimeParameters = runtimeParameters; - this.fixedParameters = fixedParameters; - this.executor = executor; - this.nodeDefinition = nodeDefinition; - this.context = context; - this.inputTypes = inputTypes; - this.outputTypes = outputTypes; - } - - @Override - public void inputReady(String inputKey, Object inputValue) { - if (parentOutputsToInputMapping.containsKey(inputKey)) { - inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue); - if (startAutomatically) start(); - } - }; - - public boolean areAllInputsReady() { - log.debug("checking inputs {} - {} ({},{})", inputs.keySet(), parentOutputsToInputMapping.keySet(), inputs.size(), parentOutputsToInputMapping.size()); - return inputs.size() == parentOutputsToInputMapping.size(); - } - - public void start() { - if (areAllInputsReady()) - new Thread(() -> { - this.context.getStepsUnderExecution().add(this.id); - log.info("starting executing of step {}", this.nodeDefinition.getName()); - try { - Map realInputs = prepareInputs(); - //log.debug("input translation for step {} is {}", id, realInputs); - Map returned = executor.execute(this.runtimeParameters, null, realInputs); - - //mapping executor output to current node output - Map mappedReturn = returned.entrySet().stream().collect(Collectors.toMap( - e -> { return !nodeDefinition.getExecutorToOutputTranslationMappings().isEmpty() - ? nodeDefinition.getExecutorToOutputTranslationMappings().get(e.getKey()) - : e.getKey(); - }, - e -> e.getValue())); - - context.addNodeResult(id, mappedReturn); - - //sendig output to next steps - nextSteps.forEach(n -> mappedReturn.forEach((k, v) -> n.inputReady(k,v))); - - } catch (ExecutionException e) { - log.error("error during execution of step {}: {}", id, e.getMessage()); - this.context.addError(this.id, e.getMessage()); - this.context.setStatus(Status.ERROR); - return; - } finally { - this.context.getStepsUnderExecution().remove(this.id); - log.info("finished executing step {}", this.nodeDefinition.getName()); - } - - }).start(); - else - log.debug("not all inputs are ready for step {}", id); - - } - - private Map prepareInputs() { - Map realInputs = new HashMap<>(); - Set keys = new HashSet<>(inputs.keySet()); - if (nodeDefinition.getInputTranslators() != null) { - //log.debug("translators {}", nodeDefinition.getInputTranslators().size()); - for (Map.Entry entry : nodeDefinition.getInputTranslators().entrySet()) { - String translation = entry.getValue().getTranslation(); - List usedKeys = entry.getValue().getUsedKeys(); - for (String usedKey : usedKeys) { - String valueToReplace = ""; - if (this.runtimeParameters.containsKey(usedKey)) - valueToReplace = (String) this.runtimeParameters.get(usedKey); - else if (inputs.containsKey(usedKey)){ - valueToReplace = (String) inputs.get(usedKey); - keys.remove(usedKey); - } - translation = translation.replace(String.format("${{%s}}", usedKey), valueToReplace); - //log.debug("replacing used key {} value {}", usedKey, valueToReplace); - } - realInputs.put(entry.getKey(), translation); - } - } - // key is the executor input name, value is the input name - if (this.nodeDefinition.getInputToExecutorTranslationMappings() != null) { - for (Map.Entry entry : this.nodeDefinition.getInputToExecutorTranslationMappings() - .entrySet()) { - if (inputs.containsKey(entry.getValue())) { - realInputs.put(entry.getKey(), inputs.get(entry.getValue())); - keys.remove(entry.getValue()); - } - } - } - - keys.forEach(k -> realInputs.put(k, inputs.get(k))); - return realInputs; } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/InputStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/InputStep.java index 94e1d97..0920ed0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/steps/InputStep.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/InputStep.java @@ -1,9 +1,10 @@ package it.cnr.isti.workflow.manager.model.steps; +import com.fasterxml.jackson.annotation.JsonTypeName; + import it.cnr.isti.workflow.manager.dto.Point; import it.cnr.isti.workflow.manager.model.types.IOType; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition; import lombok.Builder; import lombok.Data; import lombok.EqualsAndHashCode; @@ -11,6 +12,7 @@ import lombok.NonNull; @Data @EqualsAndHashCode(callSuper = true) +@JsonTypeName("input") public class InputStep extends Step implements ExecutableStep { IOType inputType; diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/InteractiveStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/InteractiveStep.java new file mode 100644 index 0000000..04819c1 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/InteractiveStep.java @@ -0,0 +1,78 @@ +package it.cnr.isti.workflow.manager.model.steps; + +import java.util.HashMap; +import java.util.Map; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonTypeName; + +import it.cnr.isti.workflow.manager.dto.Point; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition; +import it.cnr.isti.workflow.manager.model.ExecutionContext; +import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; +import it.cnr.isti.workflow.manager.model.InputReadyListener; +import it.cnr.isti.workflow.manager.model.types.IOType; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NonNull; +import lombok.ToString; + +@ToString +@Getter +@EqualsAndHashCode(callSuper = true) +@JsonTypeName("interactive") +public class InteractiveStep extends ExecutionStep implements InputReadyListener { + + private static Logger log = LoggerFactory.getLogger(InteractiveStep.class); + + IOType inputType; + IOType outputType; + + String interactionDescription; + + @JsonIgnore + ExecutionContext context; + + Map parentOutputsToInputMapping = new HashMap<>(); + + Object inputValue; + + @Builder + public InteractiveStep(@NonNull String id, @NonNull ExecutionContext context, + @NonNull IOType inputType, + @NonNull IOType outputType, @NonNull String interactionDescription, + Point position) { + super(id, position); + this.interactionDescription = interactionDescription; + this.inputType = inputType; + this.outputType = outputType; + this.context = context; + } + + @Override + public void inputReady(String key, Object value) { + log.debug("Input ready for Interactive step: {}, key: {}, value: {}", this.getId(), key, value); + this.inputValue = value; + this.context.setStatus(Status.WAITING); + this.context.getWaitingSteps().add(this.id); + } + + @Override + public void addParentOutputsToInputMapping(String parentOutputName, String inputName) { + this.parentOutputsToInputMapping.put(parentOutputName, inputName); + } + + public void setUserInput(Object input) { + this.inputValue = input; + this.context.setStatus(Status.RUNNING); + this.context.getWaitingSteps().remove(this.id); + this.getNextSteps().forEach(step -> { + step.inputReady(HumanInteractionNodeDefinition.OUTPUT_NAME, inputValue); + }); + } + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/OutputStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/OutputStep.java index 29eb7ce..228d120 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/steps/OutputStep.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/OutputStep.java @@ -7,6 +7,7 @@ import java.util.Map; import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonTypeName; import io.micrometer.common.lang.NonNull; import it.cnr.isti.workflow.manager.dto.Point; @@ -18,6 +19,7 @@ import lombok.EqualsAndHashCode; @Data @EqualsAndHashCode(callSuper = true) +@JsonTypeName("output") public class OutputStep extends Step implements InputReadyListener, ResultProvider { private IOType outputType; @@ -54,5 +56,10 @@ public class OutputStep extends Step implements InputReadyListener, ResultProvid this.resultListeners.forEach(listener -> listener.resultReady(this.id, this.getOutputName(), value)); } + @Override + public void addParentOutputsToInputMapping(String parentOutputName, String inputName) { + this.parentOutputsToInputMapping.put(parentOutputName, inputName); + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/Step.java index b54b9f2..5cd58e4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/steps/Step.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/Step.java @@ -7,7 +7,6 @@ import java.util.Map; import it.cnr.isti.workflow.manager.dto.Point; import com.fasterxml.jackson.annotation.JsonIgnore; -import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; import io.micrometer.common.lang.NonNull; @@ -18,11 +17,6 @@ import lombok.Data; @Data @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "type") -@JsonSubTypes({ - @JsonSubTypes.Type(value = InputStep.class, name = "input"), - @JsonSubTypes.Type(value = OutputStep.class, name = "output"), - @JsonSubTypes.Type(value = ExecutionStep.class, name = "execution") - }) public abstract class Step { @NonNull diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/steps/UserDefinedStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/steps/UserDefinedStep.java new file mode 100644 index 0000000..c4f106f --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/model/steps/UserDefinedStep.java @@ -0,0 +1,171 @@ +package it.cnr.isti.workflow.manager.model.steps; + +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import it.cnr.isti.workflow.manager.dto.Point; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonTypeName; + +import it.cnr.isti.workflow.manager.exceptions.ExecutionException; +import it.cnr.isti.workflow.manager.executors.Executor; +import it.cnr.isti.workflow.manager.model.ExecutionContext; +import it.cnr.isti.workflow.manager.model.InputReadyListener; +import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; +import it.cnr.isti.workflow.manager.model.types.IOType; +import it.cnr.isti.workflow.manager.model.types.Translator; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.NonNull; +import lombok.ToString; + +@ToString +@EqualsAndHashCode(callSuper = true) +@JsonTypeName("execution") +public class UserDefinedStep extends ExecutionStep implements InputReadyListener, ExecutableStep { + + private static final Logger log = org.slf4j.LoggerFactory.getLogger(UserDefinedStep.class); + + @NonNull + String name; + + Map runtimeParameters; + Map fixedParameters; + + @JsonIgnore + @NonNull + ExecutionContext context; + + Map inputTypes; + Map outputTypes; + + @JsonIgnore + Executor executor; + + Map parentOutputsToInputMapping = new HashMap<>(); + + @JsonIgnore + Map inputs = new HashMap<>(); + + @JsonIgnore + @NonNull + UserNodeDefinition nodeDefinition; + + @Builder + public UserDefinedStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map runtimeParameters, + Map fixedParameters, + UserNodeDefinition nodeDefinition, + Executor executor, Point position, @NonNull Map inputTypes, + @NonNull Map outputTypes) { + super(id, position); + this.name = name; + this.runtimeParameters = runtimeParameters; + this.fixedParameters = fixedParameters; + this.executor = executor; + this.nodeDefinition = nodeDefinition; + this.context = context; + this.inputTypes = inputTypes; + this.outputTypes = outputTypes; + } + + @Override + public void inputReady(String inputKey, Object inputValue) { + if (parentOutputsToInputMapping.containsKey(inputKey)) { + inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue); + if (startAutomatically) start(); + } + }; + + public boolean areAllInputsReady() { + log.debug("checking inputs {} - {} ({},{})", inputs.keySet(), parentOutputsToInputMapping.keySet(), inputs.size(), parentOutputsToInputMapping.size()); + return inputs.size() == parentOutputsToInputMapping.size(); + } + + public void start() { + if (areAllInputsReady()) + new Thread(() -> { + this.context.getStepsUnderExecution().add(this.id); + log.info("starting executing of step {}", this.nodeDefinition.getName()); + try { + Map realInputs = prepareInputs(); + //log.debug("input translation for step {} is {}", id, realInputs); + Map returned = executor.execute(this.runtimeParameters, null, realInputs); + + //mapping executor output to current node output + Map mappedReturn = returned.entrySet().stream().collect(Collectors.toMap( + e -> { return !nodeDefinition.getExecutorToOutputTranslationMappings().isEmpty() + ? nodeDefinition.getExecutorToOutputTranslationMappings().get(e.getKey()) + : e.getKey(); + }, + e -> e.getValue())); + + context.addNodeResult(id, mappedReturn); + + //sendig output to next steps + nextSteps.forEach(n -> mappedReturn.forEach((k, v) -> n.inputReady(k,v))); + + } catch (ExecutionException e) { + log.error("error during execution of step {}: {}", id, e.getMessage()); + this.context.addError(this.id, e.getMessage()); + this.context.setStatus(Status.ERROR); + return; + } finally { + this.context.getStepsUnderExecution().remove(this.id); + log.info("finished executing step {}", this.nodeDefinition.getName()); + } + + }).start(); + else + log.debug("not all inputs are ready for step {}", id); + + } + + private Map prepareInputs() { + Map realInputs = new HashMap<>(); + Set keys = new HashSet<>(inputs.keySet()); + if (nodeDefinition.getInputTranslators() != null) { + //log.debug("translators {}", nodeDefinition.getInputTranslators().size()); + for (Map.Entry entry : nodeDefinition.getInputTranslators().entrySet()) { + String translation = entry.getValue().getTranslation(); + List usedKeys = entry.getValue().getUsedKeys(); + for (String usedKey : usedKeys) { + String valueToReplace = ""; + if (this.runtimeParameters.containsKey(usedKey)) + valueToReplace = (String) this.runtimeParameters.get(usedKey); + else if (inputs.containsKey(usedKey)){ + valueToReplace = (String) inputs.get(usedKey); + keys.remove(usedKey); + } + translation = translation.replace(String.format("${{%s}}", usedKey), valueToReplace); + //log.debug("replacing used key {} value {}", usedKey, valueToReplace); + } + realInputs.put(entry.getKey(), translation); + } + } + // key is the executor input name, value is the input name + if (this.nodeDefinition.getInputToExecutorTranslationMappings() != null) { + for (Map.Entry entry : this.nodeDefinition.getInputToExecutorTranslationMappings() + .entrySet()) { + if (inputs.containsKey(entry.getValue())) { + realInputs.put(entry.getKey(), inputs.get(entry.getValue())); + keys.remove(entry.getValue()); + } + } + } + + keys.forEach(k -> realInputs.put(k, inputs.get(k))); + return realInputs; + } + + @Override + public void addParentOutputsToInputMapping(String parentOutputName, String inputName) { + this.parentOutputsToInputMapping.put(parentOutputName, inputName); + } + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterDefinition.java index 50803c3..39c81ab 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterDefinition.java @@ -3,7 +3,6 @@ package it.cnr.isti.workflow.manager.model.types; import java.io.IOException; import java.util.List; -import org.apache.commons.lang3.Validate; import org.json.JSONObject; import com.fasterxml.jackson.core.JacksonException; diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/InputDefinitionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/InputDefinitionRepository.java deleted file mode 100644 index 9a98069..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/InputDefinitionRepository.java +++ /dev/null @@ -1,10 +0,0 @@ -package it.cnr.isti.workflow.manager.repositories; - -import org.springframework.data.repository.CrudRepository; - -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition; - - -public interface InputDefinitionRepository extends CrudRepository { - -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/NodeDefinitionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/NodeDefinitionRepository.java index c83d059..c49eb25 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/NodeDefinitionRepository.java +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/NodeDefinitionRepository.java @@ -1,19 +1,12 @@ package it.cnr.isti.workflow.manager.repositories; -import java.util.List; - -import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; import org.springframework.stereotype.Repository; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; @Repository public interface NodeDefinitionRepository extends CrudRepository { - @Query("SELECT n.name FROM NodeDefinition n") - List getUsedTypes(); - - @Query("SELECT DISTINCT n.category FROM NodeDefinition n") - List getCategories(); + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/OutputDefinitionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/OutputDefinitionRepository.java deleted file mode 100644 index cd70bd0..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/OutputDefinitionRepository.java +++ /dev/null @@ -1,9 +0,0 @@ -package it.cnr.isti.workflow.manager.repositories; - -import org.springframework.data.repository.CrudRepository; - -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition; - -public interface OutputDefinitionRepository extends CrudRepository { - -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/SystemNodeDefinitionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/SystemNodeDefinitionRepository.java new file mode 100644 index 0000000..3d034f6 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/SystemNodeDefinitionRepository.java @@ -0,0 +1,15 @@ +package it.cnr.isti.workflow.manager.repositories; + +import java.util.List; + +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.CrudRepository; + +import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition; + + +public interface SystemNodeDefinitionRepository extends CrudRepository { + + @Query("SELECT DISTINCT n.category FROM SystemNodeDefinition n") + List getCategories(); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/UserNodeDefinitionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/UserNodeDefinitionRepository.java index 4399f4b..8910e39 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/UserNodeDefinitionRepository.java +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/UserNodeDefinitionRepository.java @@ -1,9 +1,17 @@ package it.cnr.isti.workflow.manager.repositories; +import java.util.List; + +import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; public interface UserNodeDefinitionRepository extends CrudRepository { + @Query("SELECT n.name FROM UserNodeDefinition n") + List getUsedTypes(); + + @Query("SELECT DISTINCT n.category FROM UserNodeDefinition n") + List getUsedCategories(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/services/ExecutionService.java b/src/main/java/it/cnr/isti/workflow/manager/services/ExecutionService.java index 16a0635..66dac6c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/services/ExecutionService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/services/ExecutionService.java @@ -8,7 +8,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import it.cnr.isti.workflow.manager.executors.ExecutionObject; import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; +import it.cnr.isti.workflow.manager.model.steps.ExecutionStep; import it.cnr.isti.workflow.manager.model.steps.InputStep; +import it.cnr.isti.workflow.manager.model.steps.InteractiveStep; @Service public class ExecutionService { @@ -39,45 +41,59 @@ public class ExecutionService { if (!executions.containsKey(id)) throw new IllegalArgumentException("Execution with id " + id + " not found"); if (executions.get(id).getContext().getStatus() == Status.RUNNING) - throw new IllegalStateException("Execution with id "+id+" is still running"); + throw new IllegalStateException("Execution with id " + id + " is still running"); executions.remove(id); } public ExecutionObject prepareInput(String executionId, String nodeId, String inputName, Object input) { ExecutionObject eo = getExecution(executionId); - - if (!eo.getContext().getStatus().isInitState()) { + + if (eo.getContext().getStatus().isInitState()) { + InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId)) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found")); + + if (!step.getInputName().equals(inputName)) { + throw new IllegalArgumentException("Input step with id " + nodeId + " does not accept input with name " + + inputName + " (ACCEPTED NAME is " + step.getInputName() + ")"); + } + + step.setFieldValue(input); + + if (eo.getInputSteps().stream().allMatch(s -> step.isReady())) + eo.getContext().setStatus(Status.READY); + else + eo.getContext().setStatus(Status.INITIALIZING); + } else if (eo.getContext().getStatus() == Status.WAITING) { + ExecutionStep step = eo.getExecutionSteps().stream().filter(s -> s.getId().equals(nodeId)) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found")); + + if (!(step instanceof InteractiveStep)) + throw new IllegalArgumentException("step with id " + nodeId + " is not an interactive step"); + InteractiveStep interactiveStep = (InteractiveStep) step; + if (!eo.getContext().getWaitingSteps().contains(interactiveStep.getId())) + throw new IllegalArgumentException("step with id " + nodeId + " is not in waiting state"); + + interactiveStep.setUserInput(input); + + } else throw new IllegalStateException("Execution with id " + executionId + " is not in initialization status (CURRENT STATUS is " + eo.getContext().getStatus() + ")"); - } - - InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId)) - .findFirst() - .orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found")); - - if (!step.getInputName().equals(inputName)) { - throw new IllegalArgumentException("Input step with id " + nodeId + " does not accept input with name " - + inputName + " (ACCEPTED NAME is " + step.getInputName() + ")"); - } - - step.setFieldValue(input); - - if (eo.getInputSteps().stream().allMatch(s -> step.isReady())) - eo.getContext().setStatus(Status.READY); - else - eo.getContext().setStatus(Status.INITIALIZING); - return eo; + + return eo; } public ExecutionObject startExecution(String id) { ExecutionObject eo = getExecution(id); - + if (eo.getContext().getStatus() != Status.READY) { throw new IllegalStateException("Execution with id " + id + " is not in READY status (CURRENT STATUS is " + eo.getContext().getStatus() + ")"); } - eo.getInputSteps().forEach(InputStep::start); eo.getContext().setStatus(Status.RUNNING); + eo.getInputSteps().forEach(InputStep::start); + return eo; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/services/ImportComponent.java b/src/main/java/it/cnr/isti/workflow/manager/services/ImportComponent.java index 96bcc38..f286ecd 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/services/ImportComponent.java +++ b/src/main/java/it/cnr/isti/workflow/manager/services/ImportComponent.java @@ -10,12 +10,11 @@ import org.springframework.stereotype.Component; import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.entities.FlowEntity; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; import it.cnr.isti.workflow.manager.model.auth.LoginEntity; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; import it.cnr.isti.workflow.manager.repositories.AuthRepository; import it.cnr.isti.workflow.manager.repositories.FlowRepository; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; -import jakarta.annotation.PostConstruct; @Component public class ImportComponent { diff --git a/src/main/java/it/cnr/isti/workflow/manager/services/TransformerService.java b/src/main/java/it/cnr/isti/workflow/manager/services/TransformerService.java index 7b404a0..8345d28 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/services/TransformerService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/services/TransformerService.java @@ -11,6 +11,11 @@ import org.springframework.stereotype.Service; import it.cnr.isti.workflow.manager.dto.Connection; import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.dto.Node; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.OutputNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.executors.ExecutionObject; import it.cnr.isti.workflow.manager.executors.Executor; import it.cnr.isti.workflow.manager.model.ExecutionContext; @@ -20,13 +25,11 @@ import it.cnr.isti.workflow.manager.dto.IOModel; import it.cnr.isti.workflow.manager.model.flows.errors.FlowError; import it.cnr.isti.workflow.manager.model.steps.ExecutionStep; import it.cnr.isti.workflow.manager.model.steps.InputStep; +import it.cnr.isti.workflow.manager.model.steps.InteractiveStep; import it.cnr.isti.workflow.manager.model.steps.OutputStep; import it.cnr.isti.workflow.manager.model.steps.Step; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; +import it.cnr.isti.workflow.manager.model.steps.UserDefinedStep; import it.cnr.isti.workflow.manager.model.types.IOType; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition; import static it.cnr.isti.workflow.manager.model.flows.errors.FlowError.error; import static it.cnr.isti.workflow.manager.model.flows.errors.ErrorType.*; @@ -62,20 +65,35 @@ public class TransformerService { Step step = switch (node.getNodeDefinition()) { case InputNodeDefinition id -> { - String inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY); + String inputName; + if (node.getParameters()==null) + inputName = (String) id.getRuntimeParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY).getDefaultValue(); + else inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY); + IOType type = IOType.fromString(node.getOutputs().getFirst().getType()); outputs.put(node.getOutputs().getFirst().getKey(), new IDPair(node.getKey(), inputName)); yield InputStep.builder().id(node.getKey()).position(node.getPosition()).inputName(inputName).type(type).build(); } case UserNodeDefinition ud -> { - ExecutionStep execStep = getExecutionStep(node, ud, execObject.getContext()); + UserDefinedStep execStep = getUserDefinedStep(node, ud, execObject.getContext()); + node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName()))); + node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName()))); + yield execStep; + } + case HumanInteractionNodeDefinition uid -> { + InteractiveStep execStep = getInteractiveStep(node, uid, execObject.getContext()); node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName()))); node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName()))); yield execStep; } case OutputNodeDefinition od -> { + String outputName; + if (node.getParameters()==null) + outputName = (String) od.getRuntimeParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY).getDefaultValue(); + else outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY); + IOType type = IOType.fromString(node.getInputs().getFirst().getType()); - String outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY); + inputs.put(node.getInputs().getFirst().getKey(), new IDPair(node.getKey(), outputName)); yield OutputStep.builder().id(node.getKey()).position(node.getPosition()).type(type).outputName(outputName).build(); } @@ -97,8 +115,7 @@ public class TransformerService { Step outputNode = steps.get(outputPair.nodeId); InputReadyListener inputNode = (InputReadyListener) steps.get(inputPair.nodeId); outputNode.getNextSteps().add(inputNode); - inputNode.getParentOutputsToInputMapping().put(outputPair.name, inputPair.name); - + inputNode.addParentOutputsToInputMapping(outputPair.name, inputPair.name); outputNode.addInputMapping(outputPair.name, new FieldKey(inputPair.nodeId, inputPair.name)); } @@ -120,11 +137,11 @@ public class TransformerService { return execObject; } - private ExecutionStep getExecutionStep(Node node, UserNodeDefinition nodeDefinition, ExecutionContext context) { + private UserDefinedStep getUserDefinedStep(Node node, UserNodeDefinition nodeDefinition, ExecutionContext context) { if (!executors.containsKey(nodeDefinition.getExecutor())) throw new IllegalArgumentException("Executor not found for node type: " + node.getType()); - ExecutionStep step = ExecutionStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition()) + UserDefinedStep step = UserDefinedStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition()) .inputTypes(nodeDefinition.getInputs()) .outputTypes(nodeDefinition.getOutputs()) .executor(executors.get(nodeDefinition.getExecutor())) @@ -133,6 +150,21 @@ public class TransformerService { return step; } + private InteractiveStep getInteractiveStep(Node node, HumanInteractionNodeDefinition nodeDefinition, ExecutionContext context) { + IOType inputType = IOType.fromString(node.getInputs().getFirst().getType()); + IOType outputType = IOType.fromString(node.getOutputs().getFirst().getType()); + String actionDescription = (String) node.getParameters().get(HumanInteractionNodeDefinition.PARAMETER_NAME); + InteractiveStep step = InteractiveStep.builder().id(node.getKey()).position(node.getPosition()) + .inputType(inputType) + .outputType(outputType) + .context(context) + .interactionDescription(actionDescription) + .build(); + return step; + } + + + public List isTrasformableFlow(Flow flow) { List errors = new ArrayList<>(); if (flow.getNodes().isEmpty()) diff --git a/src/test/java/it/cnr/isti/workflow/manager/Configurations.java b/src/test/java/it/cnr/isti/workflow/manager/Configurations.java index 9d62b85..33a41f8 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/Configurations.java +++ b/src/test/java/it/cnr/isti/workflow/manager/Configurations.java @@ -1,6 +1,5 @@ package it.cnr.isti.workflow.manager; -import org.aspectj.weaver.ast.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; diff --git a/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java b/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java index 1d29538..21681a6 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java @@ -9,9 +9,13 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.context.annotation.Import; import org.springframework.test.context.TestPropertySource; +import org.springframework.util.ResourceUtils; import com.fasterxml.jackson.databind.ObjectMapper; +import it.cnr.isti.workflow.manager.dto.Flow; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition; import it.cnr.isti.workflow.manager.executors.ExecutionObject; import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; import it.cnr.isti.workflow.manager.services.ExecutionService; @@ -58,4 +62,35 @@ public class ExecutionServiceTest { }); } + @Test() + void testInteractionExecutionRunningSuccess() throws Exception{ + ObjectMapper objectMapper = new ObjectMapper(); + Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:interactionFlow.json"), Flow.class); + ExecutionObject execution = executionService.createExecution(flow); + ObjectMapper mapper = new ObjectMapper(); + log.info("Execution created: {}", mapper.writeValueAsString(execution)); + execution = executionService.prepareInput(execution.getId(), "59d68aed-6de8-4305-9a3c-249069474d2e", InputNodeDefinition.DEFAULT_VALUE, "ciao como stoi"); + execution = executionService.startExecution(execution.getId()); + + while (execution.getContext().getStatus().isRunningState()) + try { + log.info("Execution status: {}", execution.getContext().getStatus()); + Thread.sleep(1000); + } catch (InterruptedException e) { } + + assert execution.getContext().getStatus() == Status.WAITING; + + String waitingStepId = execution.getContext().getWaitingSteps().getFirst(); + log.info("Waiting step id: {}", waitingStepId); + + executionService.prepareInput(execution.getId(), waitingStepId, HumanInteractionNodeDefinition.INPUT_NAME, "ciao come stai?"); + + assert execution.getContext().getResult().size() == 1; + + execution.getContext().getResult().forEach((k, v) -> { + log.info("Output: {} = {}", k, v); + }); + + } + } \ No newline at end of file diff --git a/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java b/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java index f86c7ef..71a3af9 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java @@ -88,10 +88,19 @@ public class FlowTest { void testInvalidTrasformable() throws Exception { ObjectMapper objectMapper = new ObjectMapper(); Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:invalidFlow.json"), Flow.class); - List errors = transformerService.isTrasformableFlow(flow); assertNotNull(errors); assertTrue(!errors.isEmpty()); } + @Test + void testInteractionFlowTrasformable() throws Exception { + ObjectMapper objectMapper = new ObjectMapper(); + Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:interactionFlow.json"), Flow.class); + List errors = transformerService.isTrasformableFlow(flow); + assertNotNull(errors); + log.debug("Errors: {}", errors); + assertTrue(errors.isEmpty()); + } + } diff --git a/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java b/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java index 4905ba1..f073cd2 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java +++ b/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java @@ -8,10 +8,10 @@ import it.cnr.isti.workflow.manager.dto.Connection; import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.dto.IOModel; import it.cnr.isti.workflow.manager.dto.Node; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.model.types.IOType; import it.cnr.isti.workflow.manager.model.types.Translator; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; @@ -93,8 +93,6 @@ public class TestFlowCreator { new Connection(UUID.randomUUID().toString(), execOutputKey, outputInputKey) )); - flow.getNodes().forEach(Node::resolveNodeDefinition); - return flow; } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/repositories/NodeDefRepoTest.java b/src/test/java/it/cnr/isti/workflow/manager/repositories/NodeDefRepoTest.java index 8acea5d..7265beb 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/repositories/NodeDefRepoTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/repositories/NodeDefRepoTest.java @@ -6,12 +6,12 @@ import org.springframework.test.context.TestPropertySource; import com.fasterxml.jackson.databind.ObjectMapper; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition; +import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition; import it.cnr.isti.workflow.manager.model.types.IOType; import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; import it.cnr.isti.workflow.manager.model.types.ParameterType; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; -import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition; import org.springframework.beans.factory.annotation.Autowired; @@ -83,7 +83,7 @@ class NodeDefRepoTest { private NodeDefinitionRepository nodeRepository; @Autowired - private InputDefinitionRepository inputNodeRepository; + private SystemNodeDefinitionRepository systemNodeRepository; @Test @@ -91,7 +91,7 @@ class NodeDefRepoTest { NodeDefinition node = UserNodeDefinition.builder().name("DefectDetection") .executor("GENERIC-AI").category("DefectDetection").createdBy("lucio.lelii").input("file", IOType.CSV) .fixedParameters(Map.of("LLM", - ParameterDefinition.builder().name("param").label("param").type(ParameterType.Text).build())) + ParameterDefinition.builder().name("param").label("param").defaultValue("google-gemini").type(ParameterType.Text).build())) .output("file", IOType.CSV).build(); nodeRepository.save(node); NodeDefinition foundNode = nodeRepository.findById(node.getName()).orElse(null); @@ -113,7 +113,7 @@ class NodeDefRepoTest { @Test void getSystemNodeDefinitions() { - Iterable nodes = inputNodeRepository.findAll(); + Iterable nodes = systemNodeRepository.findAll(); assertNotNull(nodes); long count = StreamSupport.stream(nodes.spliterator(), false).count(); //assertTrue(count == IOType.values().length); @@ -125,7 +125,7 @@ class NodeDefRepoTest { NodeDefinition node = UserNodeDefinition.builder().name("DefectDetection") .executor("GENERIC-AI").category("DefectDetection").createdBy("lucio.lelii").input("file", IOType.CSV) .fixedParameters(Map.of("LLM", - ParameterDefinition.builder().name("param").label("param").type(ParameterType.Text).build())) + ParameterDefinition.builder().name("param").label("param").defaultValue("google-gemini").type(ParameterType.Text).build())) .output("file", IOType.CSV).build(); nodeRepository.save(node); NodeDefinition foundNode = nodeRepository.findById(node.getName()).orElse(null); diff --git a/src/test/resources/interactionFlow.json b/src/test/resources/interactionFlow.json new file mode 100644 index 0000000..9b40c9c --- /dev/null +++ b/src/test/resources/interactionFlow.json @@ -0,0 +1,168 @@ +{ + "createdBy": "lucio.lelii", + "name": "interaction flow", + "description": null, + "nodes": [ + { + "key": "0414f102-fa57-47a5-994d-575bee0443fa", + "name": "newNode", + "createdBy": null, + "outputs": [ + { + "key": "8dda93d6-7bc6-485e-8375-7c56ca676979", + "type": "TEXT", + "name": "output" + } + ], + "inputs": [ + { + "key": "57a64d4c-bbbf-4c56-a0d2-7a2b6de2da09", + "type": "TEXT", + "name": "input" + } + ], + "color": "#A8E6CF", + "position": { + "x": 325, + "y": 145 + }, + "parameters": { + "action_description": "correct the phrase" + }, + "description": null, + "nodeDefinition": { + "type": "system", + "name": "Human Interaction", + "createdBy": "system", + "category": "humanInteraction", + "color": "black", + "fixedParameters": {}, + "runtimeParameters": { + "action_description": { + "name": "action_description", + "label": "Human Readable Action Description", + "description": "Description of the user action to perform", + "type": "Text", + "required": true, + "defaultValue": null, + "validations": [], + "specificAttributes": null + } + }, + "inputs": { + "input": "TEXT" + }, + "outputs": { + "output": "TEXT" + }, + "simulable": false + }, + "type": "Human Interaction" + }, + { + "key": "59d68aed-6de8-4305-9a3c-249069474d2e", + "name": "newNode", + "createdBy": null, + "outputs": [ + { + "key": "05d80c69-0c16-454a-ac81-ae8fa1278667", + "type": "TEXT", + "name": "value" + } + ], + "inputs": [], + "color": "#A8E6CF", + "position": { + "x": 99, + "y": 214 + }, + "parameters": null, + "description": null, + "nodeDefinition": { + "type": "system", + "name": "Text Input", + "createdBy": null, + "category": "input", + "color": "green", + "fixedParameters": {}, + "runtimeParameters": { + "name": { + "name": "name", + "label": "Input Name", + "description": "Input name", + "type": "String", + "required": false, + "defaultValue": "output", + "validations": [], + "specificAttributes": null + } + }, + "inputs": null, + "outputs": { + "value": "TEXT" + }, + "simulable": false + }, + "type": "Text Input" + }, + { + "key": "f9095513-fa24-4238-9b08-2a2b959e3978", + "name": "newNode", + "createdBy": null, + "outputs": [], + "inputs": [ + { + "key": "8daefe93-5998-4362-b24a-5087c1b2f859", + "type": "TEXT", + "name": "value" + } + ], + "color": "#A8E6CF", + "position": { + "x": 697, + "y": 200 + }, + "parameters": null, + "description": null, + "nodeDefinition": { + "type": "system", + "name": "Text Output", + "createdBy": null, + "category": "output", + "color": "red", + "fixedParameters": {}, + "runtimeParameters": { + "name": { + "name": "name", + "label": "Output Name", + "description": "Output name", + "type": "String", + "required": false, + "defaultValue": "output", + "validations": [], + "specificAttributes": null + } + }, + "inputs": { + "value": "TEXT" + }, + "outputs": null, + "simulable": false + }, + "type": "Text Output" + } + ], + "connections": [ + { + "key": "e8474793-a8d9-4fd2-9da6-1a7c5b8657d2", + "from": "05d80c69-0c16-454a-ac81-ae8fa1278667", + "to": "57a64d4c-bbbf-4c56-a0d2-7a2b6de2da09" + }, + { + "key": "7f21e4aa-d554-4fdf-ba95-c331a48ccf92", + "from": "8dda93d6-7bc6-485e-8375-7c56ca676979", + "to": "8daefe93-5998-4362-b24a-5087c1b2f859" + } + ], + "public": false +} \ No newline at end of file