From 4200a2bc3d3c65028f568c8f2da812acbedd511d Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Thu, 24 Jul 2025 15:33:39 +0200 Subject: [PATCH] refactored a lot --- .../manager/WorkflowManagerApplication.java | 11 +- .../controllers/ExecutionController.java | 14 +- .../manager/controllers/FlowsController.java | 8 +- .../{model/flows => dto}/Connection.java | 2 +- .../cnr/isti/workflow/manager/dto/Flow.java | 40 ++++++ .../manager/{model/flows => dto}/IOModel.java | 2 +- .../cnr/isti/workflow/manager/dto/Node.java | 58 ++++++++ .../cnr/isti/workflow/manager/dto/Point.java | 13 ++ .../workflow/manager/entities/FlowEntity.java | 128 ++++++++++++++++++ .../manager/executors/ExecutionObject.java | 30 ++-- .../executors/ai/services/GeminiService.java | 2 +- .../ai/services/ollama/OllamaService.java | 4 +- .../manager/model/ExecutionContext.java | 18 +-- .../isti/workflow/manager/model/FieldKey.java | 8 ++ .../workflow/manager/model/flows/Flow.java | 63 --------- .../workflow/manager/model/flows/Node.java | 50 ------- .../manager/model/steps/ExecutionStep.java | 33 +++-- .../manager/model/steps/InputStep.java | 12 +- .../manager/model/steps/OutputStep.java | 16 ++- .../workflow/manager/model/steps/Step.java | 34 +++++ .../model/types/ParameterDefinition.java | 23 +++- .../manager/model/types/ParameterType.java | 1 + .../definitions/InputNodeDefinition.java | 4 +- .../definitions/OutputNodeDefinition.java | 4 +- .../manager/repositories/FlowRepository.java | 8 +- .../converters/ConnectionListConverter.java | 2 +- .../converters/NodeListConverter.java | 5 +- .../manager/services/ExecutionService.java | 19 ++- .../manager/services/ImportComponent.java | 14 +- .../manager/services/RepositoryHolder.java | 26 ++++ .../manager/services/TransformerService.java | 54 ++++---- .../resources/workflow-editor-init/flows.json | 2 +- .../manager/ExecutionServiceTest.java | 8 +- .../cnr/isti/workflow/manager/FlowTest.java | 2 +- .../workflow/manager/TestFlowCreator.java | 13 +- .../executors/GenericAIExecutorTest.java | 1 - 36 files changed, 498 insertions(+), 234 deletions(-) rename src/main/java/it/cnr/isti/workflow/manager/{model/flows => dto}/Connection.java (84%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java rename src/main/java/it/cnr/isti/workflow/manager/{model/flows => dto}/IOModel.java (80%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/dto/Node.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/dto/Point.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/model/FieldKey.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/model/flows/Flow.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/model/flows/Node.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/services/RepositoryHolder.java 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 dcd0769..2dbcd20 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java +++ b/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java @@ -7,15 +7,15 @@ 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.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.AuthRepository; -import it.cnr.isti.workflow.manager.repositories.FlowRepository; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; +import it.cnr.isti.workflow.manager.services.ImportComponent; @SpringBootApplication public class WorkflowManagerApplication { @@ -23,13 +23,16 @@ public class WorkflowManagerApplication { @Autowired List aiModels; + @Autowired + ImportComponent importComponent; + public static void main(String[] args) { SpringApplication.run(WorkflowManagerApplication.class, args); } @Bean @ConditionalOnProperty(prefix = "app", name = "db.init.enabled", havingValue = "true") - CommandLineRunner init(NodeDefinitionRepository repository, FlowRepository flowRepository, AuthRepository authRepository) { + CommandLineRunner init(NodeDefinitionRepository repository) { return args -> { NodeDefinition textInput = InputNodeDefinition.builder().name("Text Input") @@ -41,6 +44,8 @@ public class WorkflowManagerApplication { .build(); repository.save(textInput); repository.save(textOutput); + + importComponent.start(); }; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionController.java index 2b58658..ddab9c4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionController.java @@ -18,10 +18,10 @@ import org.springframework.web.multipart.MultipartFile; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.security.SecurityRequirement; +import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.exceptions.WebException; import it.cnr.isti.workflow.manager.executors.ExecutionObject; import it.cnr.isti.workflow.manager.model.ExecutionContext; -import it.cnr.isti.workflow.manager.model.flows.Flow; import it.cnr.isti.workflow.manager.repositories.FlowRepository; import it.cnr.isti.workflow.manager.services.ExecutionService; @@ -178,11 +178,11 @@ public class ExecutionController { * @return An {@link ExecutionObject} representing the updated execution state * after preparing the input. */ - @PutMapping(path = "{id}/input/{inputName}") + @PutMapping(path = "{id}/node/{nodeId}/input/{inputName}/text", consumes = "text/plain") @Operation(summary = "Prepares string inputs", description = "Prepares the input for an execution by associating a given input string with a specific input name and execution ID.") - public ExecutionObject prepareStringInputs(@RequestBody String input, @PathVariable String inputName, + public ExecutionObject prepareStringInputs(@RequestBody String input, @PathVariable String nodeId, @PathVariable String inputName, @PathVariable String id) { - return executionService.prepareInput(id, inputName, input); + return executionService.prepareInput(id, nodeId, inputName, input); } @@ -197,14 +197,14 @@ public class ExecutionController { * preparing the input. * @throws WebServerException If an error occurs while creating or transferring the file. */ - @PutMapping(path = "{id}/input/{inputName}", consumes = "multipart/form-data") + @PutMapping(path = "{id}/node/{nodeId}/input/{inputName}/file", consumes = "multipart/form-data") @Operation(summary = "Prepares file inputs", description = "Prepares file inputs for a specific execution by uploading a file and associating it with the given input name and execution ID.") - public ExecutionObject prepareFileInputs(@RequestParam("file") MultipartFile file, @PathVariable String inputName, + public ExecutionObject prepareFileInputs(@RequestParam("file") MultipartFile file, @PathVariable String nodeId, @PathVariable String inputName, @PathVariable String id) { try { File myFile = File.createTempFile(inputName, file.getOriginalFilename()); file.transferTo(myFile); - return executionService.prepareInput(id, inputName, myFile); + return executionService.prepareInput(id, nodeId, inputName, myFile); } catch (IOException e) { throw new WebServerException("Error while creating file", e); } 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 8125e17..7488ca7 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 @@ -18,7 +18,8 @@ import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.security.SecurityRequirement; import it.cnr.isti.workflow.manager.exceptions.WebException; import it.cnr.isti.workflow.manager.model.auth.LoginEntity; -import it.cnr.isti.workflow.manager.model.flows.Flow; +import it.cnr.isti.workflow.manager.dto.Flow; +import it.cnr.isti.workflow.manager.entities.FlowEntity; import it.cnr.isti.workflow.manager.model.flows.errors.FlowError; import it.cnr.isti.workflow.manager.repositories.FlowRepository; import it.cnr.isti.workflow.manager.services.TransformerService; @@ -27,7 +28,6 @@ import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.PutMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.GetMapping; -import org.springframework.web.bind.annotation.RequestParam; @@ -90,7 +90,7 @@ 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(flow); + Flow savedFlow = flowRepository.save((FlowEntity)flow); logger.info("flow created {}", flow.getId()); return savedFlow; @@ -118,7 +118,7 @@ public class FlowsController { if (!savedFlow.getCreatedBy().equals(userDetails.getUsername())) throw new WebException(HttpStatus.FORBIDDEN, "Flow not created by user"); - savedFlow = flowRepository.save(flow); + savedFlow = flowRepository.save(FlowEntity.from(flow)); logger.info("Flow saved: {} ",savedFlow.getId()); return savedFlow; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Connection.java b/src/main/java/it/cnr/isti/workflow/manager/dto/Connection.java similarity index 84% rename from src/main/java/it/cnr/isti/workflow/manager/model/flows/Connection.java rename to src/main/java/it/cnr/isti/workflow/manager/dto/Connection.java index 9927993..858bdb1 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Connection.java +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Connection.java @@ -1,4 +1,4 @@ -package it.cnr.isti.workflow.manager.model.flows; +package it.cnr.isti.workflow.manager.dto; import lombok.AllArgsConstructor; import lombok.Builder; 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 new file mode 100644 index 0000000..c57e5ea --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Flow.java @@ -0,0 +1,40 @@ +package it.cnr.isti.workflow.manager.dto; + +import java.util.List; +import com.fasterxml.jackson.annotation.JsonAlias; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Builder.Default; +import lombok.Data; +import lombok.NoArgsConstructor; +import lombok.NonNull; +import lombok.ToString; + +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +@ToString +public class Flow { + + String id; + + @NonNull + String createdBy; + + @Default + @JsonAlias("public") + boolean isPublic =false; + + @NonNull + String name; + + String description; + + @Builder.Default + List nodes = List.of(); + + @Builder.Default + List connections = List.of(); + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/flows/IOModel.java b/src/main/java/it/cnr/isti/workflow/manager/dto/IOModel.java similarity index 80% rename from src/main/java/it/cnr/isti/workflow/manager/model/flows/IOModel.java rename to src/main/java/it/cnr/isti/workflow/manager/dto/IOModel.java index bcfbcaf..6631501 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/flows/IOModel.java +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/IOModel.java @@ -1,4 +1,4 @@ -package it.cnr.isti.workflow.manager.model.flows; +package it.cnr.isti.workflow.manager.dto; import lombok.AllArgsConstructor; import lombok.Data; 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 new file mode 100644 index 0000000..e69aab0 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Node.java @@ -0,0 +1,58 @@ +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.services.RepositoryHolder; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; +import lombok.Singular; +import lombok.ToString; + +@Data +@Builder +@AllArgsConstructor +@NoArgsConstructor +@ToString +public class Node { + + String key; + String name; + + String createdBy; + + @Singular + List outputs; + + @Singular + List inputs; + + String color; + Point position; + + Map parameters; + + String description; + + @JsonProperty(access = JsonProperty.Access.READ_ONLY) + NodeDefinition nodeDefinition; + + String type; + + public void resolveNodeDefinition() { + if (type != null) { + this.nodeDefinition = RepositoryHolder.getNodeDefinition() + .findById(type) + .orElse(null); + } + } +} + + + diff --git a/src/main/java/it/cnr/isti/workflow/manager/dto/Point.java b/src/main/java/it/cnr/isti/workflow/manager/dto/Point.java new file mode 100644 index 0000000..5e84702 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/dto/Point.java @@ -0,0 +1,13 @@ +package it.cnr.isti.workflow.manager.dto; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@NoArgsConstructor +@AllArgsConstructor +@Data +public class Point { + int x; + int y; +} 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 new file mode 100644 index 0000000..bbe5bd2 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/entities/FlowEntity.java @@ -0,0 +1,128 @@ +package it.cnr.isti.workflow.manager.entities; + + +import java.util.List; + +import org.hibernate.annotations.UuidGenerator; +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.repositories.converters.ConnectionListConverter; +import it.cnr.isti.workflow.manager.repositories.converters.NodeListConverter; +import jakarta.persistence.Column; +import jakarta.persistence.Convert; +import jakarta.persistence.Entity; +import jakarta.persistence.Id; +import jakarta.persistence.PostLoad; +import lombok.EqualsAndHashCode; + +@Entity +@EqualsAndHashCode(callSuper = true) +public class FlowEntity extends Flow { + + public FlowEntity() { + super(); + } + + public FlowEntity(Flow flow) { + super( + flow.getId(), + flow.getCreatedBy(), + flow.isPublic(), + flow.getName(), + flow.getDescription(), + flow.getNodes(), + flow.getConnections() + ); + } + + public static FlowEntity from(Flow flow) { + return new FlowEntity(flow); + } + + + @Id + @UuidGenerator + @Override + public String getId() { + return super.getId(); + } + + @Override + public void setId(String id) { + super.setId(id); + } + + @Column(nullable = false) + @Override + public String getCreatedBy() { + return super.getCreatedBy(); + } + + @Override + public void setCreatedBy(String createdBy) { + super.setCreatedBy(createdBy); + } + + @Column(name = "public", nullable = false) + @Override + public boolean isPublic() { + return super.isPublic(); + } + + @Override + public void setPublic(boolean isPublic) { + super.setPublic(isPublic); + } + + @Column(nullable = false) + @Override + public String getName() { + return super.getName(); + } + + @Override + public void setName(String name) { + super.setName(name); + } + + @Column + @Override + public String getDescription() { + return super.getDescription(); + } + + @Override + public void setDescription(String description) { + super.setDescription(description); + } + + @Convert(converter = NodeListConverter.class) + @Column(columnDefinition = "TEXT") + @Override + public List getNodes() { + return super.getNodes(); + } + + @Override + public void setNodes(List nodes) { + super.setNodes(nodes); + } + + @Convert(converter = ConnectionListConverter.class) + @Column(columnDefinition = "TEXT") + @Override + public List getConnections() { + return super.getConnections(); + } + + @Override + public void setConnections(List connections) { + super.setConnections(connections); + } + + @PostLoad + public void postLoad() { + this.getNodes().forEach(Node::resolveNodeDefinition); + } +} 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 edae237..d166623 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 @@ -7,9 +7,11 @@ 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; -import it.cnr.isti.workflow.manager.model.flows.Flow; +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.OutputStep; import it.cnr.isti.workflow.manager.model.steps.ResultListener; @@ -30,30 +32,28 @@ public class ExecutionObject implements ResultListener { @Setter String name; - Map exepectedOutputsDescription; + Map exepectedOutputsDescription; - Map exepectedInputsDescription; + Map exepectedInputsDescription; - Flow flow; - @JsonIgnore List inputSteps; - @JsonIgnore List outputSteps; + List executionSteps; - public ExecutionObject(Flow flow) { - this.flow = flow; - this.name = flow.getName()+" - "+this.creationTime; + + public ExecutionObject(String executionName) { + this.name = executionName; } public void setOutputSteps(List outputSteps) { this.outputSteps = outputSteps; this.exepectedOutputsDescription = new HashMap<>(); for (OutputStep outputStep : outputSteps) { - this.exepectedOutputsDescription.put(outputStep.getOutputName(), outputStep.getType()); + this.exepectedOutputsDescription.put(new FieldKey(outputStep.getId(), outputStep.getOutputName()), outputStep.getOutputType()); outputStep.registerResultListener(this); } } @@ -61,8 +61,12 @@ public class ExecutionObject implements ResultListener { public void setInputSteps(List inputSteps) { this.inputSteps = inputSteps; this.exepectedInputsDescription = new HashMap<>(); - for (InputStep inputStep : inputSteps) - this.exepectedInputsDescription.put(inputStep.getInputName(), inputStep.getType()); + for (InputStep inputStep : inputSteps) + this.exepectedInputsDescription.put(new FieldKey(inputStep.getId(), inputStep.getInputName()), inputStep.getInputType()); + } + + public void setExecutionSteps(List executionSteps) { + this.executionSteps = executionSteps; } @Override @@ -70,7 +74,7 @@ public class ExecutionObject implements ResultListener { //TODO: devo discriminare in base al nodeId altrimenti se ho più output step con lo stesso inputName non so a quale output step riferirmi - this.context.addResult(resultKey, value); + this.context.addResult(nodeKey, resultKey, value); if (this.exepectedOutputsDescription.keySet().containsAll(this.context.getExecutionResult().keySet())) { //significa che l'esecuzione è finita this.context.setStatus(Status.SUCCESS); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/GeminiService.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/GeminiService.java index 2da3633..996d6c2 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/GeminiService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/GeminiService.java @@ -47,7 +47,7 @@ public class GeminiService { status -> status.is5xxServerError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("Error: server without body") - .flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body)))) + .flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body)))) .bodyToMono(String.class) .timeout(Duration.ofMinutes(2)) .retryWhen( diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/ollama/OllamaService.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/ollama/OllamaService.java index 3e6917e..48690b0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/ollama/OllamaService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/services/ollama/OllamaService.java @@ -67,7 +67,7 @@ public class OllamaService { status -> status.is5xxServerError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("Error: server without body") - .flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body)))) + .flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body)))) .bodyToMono(GenerateResponse.class) // deserialize JSON in oggetto Java .timeout(Duration.ofMinutes(2)); @@ -86,7 +86,7 @@ public class OllamaService { status -> status.is5xxServerError(), clientResponse -> clientResponse.bodyToMono(String.class) .defaultIfEmpty("Error: server without body") - .flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body)))) + .flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body)))) .bodyToMono(ModelResponse.class) // deserialize JSON in oggetto Java .timeout(Duration.ofMinutes(1)) .map(modelResponse -> modelResponse.getModels().stream() 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 15d3b23..8fd89e4 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 @@ -5,8 +5,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; - -import lombok.AccessLevel; import lombok.Getter; import lombok.NoArgsConstructor; @@ -47,14 +45,14 @@ public class ExecutionContext { Map inputs = new HashMap<>(); @Getter() - private Map result = new HashMap<>(); + private Map result = new HashMap<>(); Long startTime = null; Long endTime = null; List stepsUnderExecution = new ArrayList<>(); - List errors = new ArrayList<>();; - List warnings = new ArrayList<>();; + Map errors = new HashMap<>(); + List warnings = new ArrayList<>(); Map> nodeResult = new HashMap<>(); @@ -74,15 +72,19 @@ public class ExecutionContext { this.nodeResult.put(nodeId, result); } - public void addResult(String key, Object value) { - this.result.put(key, value); + public void addResult(String nodeId, String key, Object value) { + this.result.put(new FieldKey(nodeId, key), value); } public void addInput(String key, Object value) { this.inputs.put(key, value); } - public Map getExecutionResult() { + public void addError(String nodeId, String error) { + this.errors.put(nodeId, error); + } + + public Map getExecutionResult() { return Collections.unmodifiableMap(this.result); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/FieldKey.java b/src/main/java/it/cnr/isti/workflow/manager/model/FieldKey.java new file mode 100644 index 0000000..6859873 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/model/FieldKey.java @@ -0,0 +1,8 @@ +package it.cnr.isti.workflow.manager.model; + +public record FieldKey(String nodeId, String fieldId) { + @Override + public final String toString() { + return nodeId + ":" + fieldId; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Flow.java b/src/main/java/it/cnr/isti/workflow/manager/model/flows/Flow.java deleted file mode 100644 index d7ffdc1..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Flow.java +++ /dev/null @@ -1,63 +0,0 @@ -package it.cnr.isti.workflow.manager.model.flows; - - -import java.util.List; - -import org.hibernate.annotations.UuidGenerator; - -import com.fasterxml.jackson.annotation.JsonAlias; - -import it.cnr.isti.workflow.manager.repositories.converters.ConnectionListConverter; -import it.cnr.isti.workflow.manager.repositories.converters.NodeListConverter; -import jakarta.persistence.Column; -import jakarta.persistence.Convert; -import jakarta.persistence.Entity; -import jakarta.persistence.Id; -import jakarta.validation.constraints.NotBlank; -import lombok.AllArgsConstructor; -import lombok.Builder; -import lombok.Builder.Default; -import lombok.Data; -import lombok.NoArgsConstructor; -import lombok.NonNull; -import lombok.ToString; - -@Data -@AllArgsConstructor -@NoArgsConstructor -@Builder -@ToString -@Entity -public class Flow { - - @Id - @UuidGenerator - String id; - - @NotBlank - @NonNull - @Column(name = "created_by") - String createdBy; - - @Default - @JsonAlias("public") - boolean isPublic =false; - - @NotBlank - @NonNull - String name; - - String description; - - @Builder.Default - @Convert(converter = NodeListConverter.class) - @Column(columnDefinition = "TEXT") - List nodes = List.of(); - - @Builder.Default - @Convert(converter = ConnectionListConverter.class) - @Column(columnDefinition = "TEXT") - List connections = List.of(); - - -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Node.java b/src/main/java/it/cnr/isti/workflow/manager/model/flows/Node.java deleted file mode 100644 index ab6d70b..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/model/flows/Node.java +++ /dev/null @@ -1,50 +0,0 @@ -package it.cnr.isti.workflow.manager.model.flows; - -import java.util.List; -import java.util.Map; -import lombok.AllArgsConstructor; -import lombok.Builder; -import lombok.Data; -import lombok.NoArgsConstructor; -import lombok.Singular; -import lombok.ToString; - -@Data -@Builder -@AllArgsConstructor -@NoArgsConstructor -@ToString -public class Node { - - String key; - String name; - - String createdBy; - - @Singular - List outputs; - - @Singular - List inputs; - - String color; - Point position; - - Map parameters; - - String description; - - String type; - -} - - -@NoArgsConstructor -@AllArgsConstructor -@Data -class Point { - int x; - int y; -} - - 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 de78162..c198029 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 @@ -6,17 +6,17 @@ 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.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.EqualsAndHashCode; @@ -24,40 +24,53 @@ import lombok.NonNull; import lombok.ToString; @Data -@AllArgsConstructor @ToString @EqualsAndHashCode(callSuper = true) public class ExecutionStep extends Step implements InputReadyListener, ExecutableStep { 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(ExecutionContext context, String id, Map runtimeParameters, + public ExecutionStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map runtimeParameters, Map fixedParameters, UserNodeDefinition nodeDefinition, - Executor executor) { - this.id = id; + 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 @@ -98,12 +111,12 @@ public class ExecutionStep extends Step implements InputReadyListener, Executabl } catch (ExecutionException e) { log.error("error during execution of step {}: {}", id, e.getMessage()); - this.context.getErrors().add(String.format("[%S] %s", this.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 of step {}", this.nodeDefinition.getName()); + log.info("finished executing step {}", this.nodeDefinition.getName()); } }).start(); 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 6fdb7ea..94e1d97 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,5 +1,7 @@ package it.cnr.isti.workflow.manager.model.steps; + +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; @@ -11,17 +13,17 @@ import lombok.NonNull; @EqualsAndHashCode(callSuper = true) public class InputStep extends Step implements ExecutableStep { - IOType type; + IOType inputType; Object fieldValue; String inputName; @Builder - InputStep(@NonNull String id,@NonNull String inputName, @NonNull IOType type) { - this.id = id; + InputStep(@NonNull String id,@NonNull String inputName, @NonNull IOType type, Point position) { + super(id, position); this.inputName = inputName; - this.type = type; + this.inputType = type; } public void setFieldValue(Object fieldValue) { @@ -35,7 +37,7 @@ public class InputStep extends Step implements ExecutableStep { @Override public void start() { this.getNextSteps().forEach(step -> { - step.inputReady(InputNodeDefinition.OUTPUT_NAME, this.fieldValue); + step.inputReady(this.inputName, this.fieldValue); }); } 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 dc13856..29eb7ce 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 @@ -5,6 +5,11 @@ import java.util.HashMap; import java.util.List; import java.util.Map; + +import com.fasterxml.jackson.annotation.JsonIgnore; + +import io.micrometer.common.lang.NonNull; +import it.cnr.isti.workflow.manager.dto.Point; import it.cnr.isti.workflow.manager.model.InputReadyListener; import it.cnr.isti.workflow.manager.model.types.IOType; import lombok.Builder; @@ -15,21 +20,20 @@ import lombok.EqualsAndHashCode; @EqualsAndHashCode(callSuper = true) public class OutputStep extends Step implements InputReadyListener, ResultProvider { - private String id; - - private IOType type; + private IOType outputType; private String outputName; + @JsonIgnore List resultListeners = new ArrayList<>(); Map parentOutputsToInputMapping = new HashMap<>(); @Builder - OutputStep(String id, IOType type, String outputName) { + OutputStep(@NonNull String id, @NonNull IOType type, @NonNull String outputName, Point position) { + super(id, position); this.outputName = outputName; - this.id = id; - this.type = type; + this.outputType = type; } @Override 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 f4c762d..b54b9f2 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 @@ -1,16 +1,50 @@ package it.cnr.isti.workflow.manager.model.steps; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +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; +import it.cnr.isti.workflow.manager.model.FieldKey; import it.cnr.isti.workflow.manager.model.InputReadyListener; 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 String id; + + //output name to external input for connections + Map inputMapping = new HashMap<>(); + + Point position; + + @JsonIgnore List nextSteps = new ArrayList<>(); + boolean startAutomatically = true; + public Step(String id, Point position){ + this.id = id; + this.position = position; + } + + public void addInputMapping(String outputName, FieldKey fieldKey) { + this.inputMapping.put(outputName, fieldKey); + } + } 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 dc89061..50803c3 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 @@ -2,6 +2,8 @@ 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; @@ -38,8 +40,9 @@ public class ParameterDefinition { @NonNull private ParameterType type; - @Builder.Default - private boolean required = false; + protected boolean required = false; + + private Object defaultValue; @Singular private List validations; @@ -47,6 +50,22 @@ public class ParameterDefinition { @JsonDeserialize(using = JsonObjectDeserializer.class) @JsonSerialize(using = JsonObjectSerializer.class) private JSONObject specificAttributes; + + public static ParameterDefinitionBuilder builder() { + return new CustomParameterDefinitionBuilder(); + } + + // extend the generated builder class to add custom logic + private static class CustomParameterDefinitionBuilder extends ParameterDefinitionBuilder { + + @Override + public ParameterDefinition build() { + if (!super.required && super.defaultValue == null) { + throw new IllegalArgumentException("When not required the parameter must have a default value."); + } + return super.build(); + } + } } class JsonObjectDeserializer extends JsonDeserializer { diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterType.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterType.java index 9a8d6ca..27175d3 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/ParameterType.java @@ -5,6 +5,7 @@ public enum ParameterType { Select, DynamicMultiOptions, Text, + String, Boolean, Number } 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/model/types/nodes/definitions/InputNodeDefinition.java index 9070e02..16edfe3 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/model/types/nodes/definitions/InputNodeDefinition.java @@ -30,9 +30,9 @@ public class InputNodeDefinition extends NodeDefinition { this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder() .name(INPUT_PARAMETER_KEY) .label("Input Name") - .required(true) .description("Input name") - .type(ParameterType.Text) + .defaultValue("output") + .type(ParameterType.String) .build()); } 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/model/types/nodes/definitions/OutputNodeDefinition.java index 3e76cc8..ea26521 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/model/types/nodes/definitions/OutputNodeDefinition.java @@ -32,9 +32,9 @@ public class OutputNodeDefinition extends NodeDefinition { this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder() .name(INPUT_PARAMETER_KEY) .label("Output Name") - .required(true) + .defaultValue("output") .description("Output name") - .type(ParameterType.Text) + .type(ParameterType.String) .build()); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/FlowRepository.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/FlowRepository.java index b30fe8b..6eb2f70 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/FlowRepository.java +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/FlowRepository.java @@ -7,11 +7,13 @@ import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; import org.springframework.stereotype.Repository; -import it.cnr.isti.workflow.manager.model.flows.Flow; +import it.cnr.isti.workflow.manager.dto.Flow; +import it.cnr.isti.workflow.manager.entities.FlowEntity; + @Repository -public interface FlowRepository extends JpaRepository { +public interface FlowRepository extends JpaRepository { - @Query("SELECT f FROM Flow f WHERE f.createdBy = :createdBy OR f.isPublic = true") + @Query("SELECT f FROM FlowEntity f WHERE f.createdBy = :createdBy OR f.public = true") List findFlowsByCreatedByOrPublic(@Param("createdBy") String createdBy); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/ConnectionListConverter.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/ConnectionListConverter.java index abf662d..ef83097 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/ConnectionListConverter.java +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/ConnectionListConverter.java @@ -6,7 +6,7 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; -import it.cnr.isti.workflow.manager.model.flows.Connection; +import it.cnr.isti.workflow.manager.dto.Connection; import jakarta.persistence.AttributeConverter; import jakarta.persistence.Converter; diff --git a/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/NodeListConverter.java b/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/NodeListConverter.java index 5f0e40b..5a82421 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/NodeListConverter.java +++ b/src/main/java/it/cnr/isti/workflow/manager/repositories/converters/NodeListConverter.java @@ -4,9 +4,8 @@ import java.util.List; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; -import com.fasterxml.jackson.databind.ObjectMapper; -import it.cnr.isti.workflow.manager.model.flows.Node; +import it.cnr.isti.workflow.manager.dto.Node; import it.cnr.isti.workflow.manager.services.ObjectMapperHolder; import jakarta.persistence.AttributeConverter; import jakarta.persistence.Converter; @@ -14,7 +13,7 @@ import jakarta.persistence.Converter; @Converter public class NodeListConverter implements AttributeConverter, String> { - + @Override public String convertToDatabaseColumn(List nodes) { 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 7e7852b..16a0635 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,6 @@ 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.flows.Flow; import it.cnr.isti.workflow.manager.model.steps.InputStep; @Service @@ -19,7 +18,7 @@ public class ExecutionService { @Autowired TransformerService transformerService; - public ExecutionObject createExecution(Flow flow) { + public ExecutionObject createExecution(it.cnr.isti.workflow.manager.dto.Flow flow) { ExecutionObject execObject = transformerService.transform(flow); executions.put(execObject.getId(), execObject); return execObject; @@ -44,17 +43,23 @@ public class ExecutionService { executions.remove(id); } - public ExecutionObject prepareInput(String id, String key, Object input) { - ExecutionObject eo = getExecution(id); + public ExecutionObject prepareInput(String executionId, String nodeId, String inputName, Object input) { + ExecutionObject eo = getExecution(executionId); if (!eo.getContext().getStatus().isInitState()) { - throw new IllegalStateException("Execution with id " + id + 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(key)) + InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId)) .findFirst() - .orElseThrow(() -> new IllegalArgumentException("Input step with id " + key + " not found")); + .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())) 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 7827416..96bcc38 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 @@ -8,9 +8,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; 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.model.auth.LoginEntity; -import it.cnr.isti.workflow.manager.model.flows.Flow; 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; @@ -32,6 +32,9 @@ public class ImportComponent { @Autowired NodeDefinitionRepository nodeRepository; + @Autowired + RepositoryHolder repositoryHolder; + @Autowired FlowRepository flowRepository; @@ -44,8 +47,7 @@ public class ImportComponent { @Value("${app.import.path}") String path; - @PostConstruct - void init() { + public void start() { if (!enabled) { System.out.println("ImportComponent is disabled"); return; @@ -85,7 +87,7 @@ public class ImportComponent { logger.debug("Node {} not found, saving new node", node.getName()); nodeRepository.save(node); } - ); + ); } catch (Exception e) { logger.error("Error loading nodes: {}", e.getMessage(),e); } @@ -102,7 +104,7 @@ public class ImportComponent { logger.info("Loading flows from file: {}", file.getAbsolutePath()); List flows = ObjectMapperHolder.mapper.readValue(file, ObjectMapperHolder.mapper.getTypeFactory().constructCollectionType(List.class,Flow.class)); for (Flow flow: flows){ - flowRepository.saveAndFlush(flow); + flowRepository.saveAndFlush(FlowEntity.from(flow)); } } catch (Exception e) { logger.error("Error loading flows: {}", e.getMessage(),e); diff --git a/src/main/java/it/cnr/isti/workflow/manager/services/RepositoryHolder.java b/src/main/java/it/cnr/isti/workflow/manager/services/RepositoryHolder.java new file mode 100644 index 0000000..a6abbde --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/services/RepositoryHolder.java @@ -0,0 +1,26 @@ +package it.cnr.isti.workflow.manager.services; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; +import jakarta.annotation.PostConstruct; + +@Component +public class RepositoryHolder { + + @Autowired + private NodeDefinitionRepository nodeDefinitionRepository; + + private static NodeDefinitionRepository _nodeRepository; + + @PostConstruct + public void init() { + System.out.println("Initializing RepositoryHolder with NodeDefinitionRepository"); + RepositoryHolder._nodeRepository = nodeDefinitionRepository; + } + + public static NodeDefinitionRepository getNodeDefinition() { + return _nodeRepository; + } + +} 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 4c7e7b4..7b404a0 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 @@ -8,21 +8,21 @@ import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; 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.executors.ExecutionObject; import it.cnr.isti.workflow.manager.executors.Executor; import it.cnr.isti.workflow.manager.model.ExecutionContext; +import it.cnr.isti.workflow.manager.model.FieldKey; import it.cnr.isti.workflow.manager.model.InputReadyListener; -import it.cnr.isti.workflow.manager.model.flows.Connection; -import it.cnr.isti.workflow.manager.model.flows.Flow; -import it.cnr.isti.workflow.manager.model.flows.IOModel; -import it.cnr.isti.workflow.manager.model.flows.Node; +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.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.repositories.NodeDefinitionRepository; 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; @@ -40,15 +40,8 @@ public class TransformerService { @Autowired Map executors; - @Autowired - NodeDefinitionRepository nodeDefinitionRepository; - private NodeDefinition getNodeDefinition(String nodeType) { - return nodeDefinitionRepository.findById(nodeType) - .orElseThrow(() -> new IllegalArgumentException("Node definition not found for type: " + nodeType)); - } - - public ExecutionObject transform(Flow flow) { + public ExecutionObject transform(it.cnr.isti.workflow.manager.dto.Flow flow) { if (flow == null ) throw new IllegalArgumentException("Flow cannot be null"); @@ -63,29 +56,33 @@ public class TransformerService { Map inputs = new HashMap<>(); Map outputs = new HashMap<>(); - ExecutionObject execObject = new ExecutionObject(flow); + ExecutionObject execObject = new ExecutionObject(flow.getName()); for (Node node : flow.getNodes()) { - Step step = switch (getNodeDefinition(node.getType())) { + Step step = switch (node.getNodeDefinition()) { case InputNodeDefinition id -> { String inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY); IOType type = IOType.fromString(node.getOutputs().getFirst().getType()); - yield InputStep.builder().id(node.getKey()).inputName(inputName).type(type).build(); + 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()); + 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 UserNodeDefinition ud -> getExecutionStep(node, ud, execObject.getContext()); case OutputNodeDefinition od -> { IOType type = IOType.fromString(node.getInputs().getFirst().getType()); String outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY); - yield OutputStep.builder().id(node.getKey()).type(type).outputName(outputName).build(); + 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(); } default -> throw new IllegalArgumentException("Node type " + node.getType() + " not supported"); }; - steps.put(node.getKey(), step); - 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()))); - + } for (Connection connection : flow.getConnections()) { @@ -101,6 +98,8 @@ public class TransformerService { InputReadyListener inputNode = (InputReadyListener) steps.get(inputPair.nodeId); outputNode.getNextSteps().add(inputNode); inputNode.getParentOutputsToInputMapping().put(outputPair.name, inputPair.name); + + outputNode.addInputMapping(outputPair.name, new FieldKey(inputPair.nodeId, inputPair.name)); } execObject.setInputSteps(steps.values().stream() @@ -113,6 +112,11 @@ public class TransformerService { .map(s -> (OutputStep) s) .toList()); + execObject.setExecutionSteps(steps.values().stream() + .filter(s -> s instanceof ExecutionStep) + .map(s -> (ExecutionStep) s) + .toList()); + return execObject; } @@ -120,7 +124,9 @@ public class TransformerService { if (!executors.containsKey(nodeDefinition.getExecutor())) throw new IllegalArgumentException("Executor not found for node type: " + node.getType()); - ExecutionStep step = ExecutionStep.builder().id(node.getKey()) + ExecutionStep step = ExecutionStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition()) + .inputTypes(nodeDefinition.getInputs()) + .outputTypes(nodeDefinition.getOutputs()) .executor(executors.get(nodeDefinition.getExecutor())) .runtimeParameters(node.getParameters()).nodeDefinition(nodeDefinition).context(context).build(); @@ -193,7 +199,7 @@ public class TransformerService { // CHECK PARAMETERS // TODO: devo controllare che i parametri siano del tipo giusto - NodeDefinition nodeDef = getNodeDefinition(n.getType()); + NodeDefinition nodeDef = n.getNodeDefinition(); if (nodeDef.getRuntimeParameters() !=null && n.getParameters() != null) { nodeDef.getRuntimeParameters().forEach((k, v) -> { diff --git a/src/main/resources/workflow-editor-init/flows.json b/src/main/resources/workflow-editor-init/flows.json index 8450c9c..c22e677 100644 --- a/src/main/resources/workflow-editor-init/flows.json +++ b/src/main/resources/workflow-editor-init/flows.json @@ -157,7 +157,7 @@ "y": -16 }, "parameters": { - "name": "interview\n" + "name": "interview" }, "description": null, "type": "Text Input" 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 e5382ad..1d29538 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/ExecutionServiceTest.java @@ -10,6 +10,8 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.context.annotation.Import; import org.springframework.test.context.TestPropertySource; +import com.fasterxml.jackson.databind.ObjectMapper; + 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; @@ -34,9 +36,11 @@ public class ExecutionServiceTest { } @Test() - void testExecutionRunningSuccess() { + void testExecutionRunningSuccess() throws Exception{ ExecutionObject execution = executionService.createExecution(testFlowCreator.create()); - execution = executionService.prepareInput(execution.getId(), TestFlowCreator.INPUT_KEY, "ciao come stai"); + ObjectMapper mapper = new ObjectMapper(); + log.info("Execution created: {}", mapper.writeValueAsString(execution)); + execution = executionService.prepareInput(execution.getId(), TestFlowCreator.INPUT_KEY, "input", "ciao come stai"); execution = executionService.startExecution(execution.getId()); while (execution.getContext().getStatus().isRunningState()) 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 a274a7f..f86c7ef 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/FlowTest.java @@ -15,7 +15,7 @@ import org.springframework.util.ResourceUtils; import com.fasterxml.jackson.databind.ObjectMapper; import it.cnr.isti.workflow.manager.executors.ExecutionObject; import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; -import it.cnr.isti.workflow.manager.model.flows.Flow; +import it.cnr.isti.workflow.manager.dto.Flow; import it.cnr.isti.workflow.manager.model.flows.errors.FlowError; import it.cnr.isti.workflow.manager.model.steps.ResultListener; import it.cnr.isti.workflow.manager.repositories.FlowRepository; 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 3aa1d2f..4905ba1 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java +++ b/src/test/java/it/cnr/isti/workflow/manager/TestFlowCreator.java @@ -4,10 +4,10 @@ import java.util.List; import java.util.Map; import java.util.UUID; -import it.cnr.isti.workflow.manager.model.flows.Connection; -import it.cnr.isti.workflow.manager.model.flows.Flow; -import it.cnr.isti.workflow.manager.model.flows.IOModel; -import it.cnr.isti.workflow.manager.model.flows.Node; +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.model.types.IOType; import it.cnr.isti.workflow.manager.model.types.Translator; import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition; @@ -26,7 +26,7 @@ public class TestFlowCreator { } public Flow create(){ - Translator inputFieldtransaltor = Translator.builder().translation("translate the following phrase : ${{phrase}} in french, return only the translation").build(); + Translator inputFieldtransaltor = Translator.builder().translation("translate the following phrase : ${{phrase}} in french, return only the translation").build(); NodeDefinition nFrenchTranslator = UserNodeDefinition.builder().name("FrenchTranslator") .executor("GENERIC-AI").category("Translators").createdBy("lucio.lelii").input("phrase", IOType.TEXT) @@ -53,6 +53,7 @@ public class TestFlowCreator { .input(new IOModel(execInputKey, IOType.TEXT.toString(), "phrase")) .output(new IOModel(execOutputKey, IOType.TEXT.toString(), "translated")) .parameters(parameter) + .nodeDefinition(nFrenchTranslator) .createdBy("lucio.lelii") .color("#FF0000") .build(); @@ -92,6 +93,8 @@ 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/executors/GenericAIExecutorTest.java b/src/test/java/it/cnr/isti/workflow/manager/executors/GenericAIExecutorTest.java index 3228455..7428377 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executors/GenericAIExecutorTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executors/GenericAIExecutorTest.java @@ -1,6 +1,5 @@ package it.cnr.isti.workflow.manager.executors; -import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue;