diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..dac3b4c --- /dev/null +++ b/Dockerfile @@ -0,0 +1,4 @@ +FROM eclipse-temurin:17-jdk-alpine +VOLUME /tmp +COPY target/workflow-manager.jar app.jar +ENTRYPOINT ["java", "-jar", "/workflow-manager.jar"] \ No newline at end of file diff --git a/pom.xml b/pom.xml index 93c5769..339013a 100644 --- a/pom.xml +++ b/pom.xml @@ -12,7 +12,7 @@ it.cnr.isti workflow-manager 0.0.1-SNAPSHOT - war + jar workflow-manager Workflow server project @@ -49,17 +49,17 @@ org.springframework.boot spring-boot-starter-webflux + + org.springdoc + springdoc-openapi-starter-webmvc-ui + 2.0.2 + org.projectlombok lombok true - - org.springframework.boot - spring-boot-starter-tomcat - provided - org.json @@ -75,6 +75,7 @@ + workflow-manager org.apache.maven.plugins diff --git a/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java b/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java new file mode 100644 index 0000000..1e4d36c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/ExecutionService.java @@ -0,0 +1,65 @@ +package it.cnr.isti.workflow.manager; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +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.ui.FlowEntity; + +@Service +public class ExecutionService { + + private static Map executions = new HashMap<>(); + + @Autowired + TransformerService transformerService; + + + public ExecutionObject createExecution(FlowEntity flow) { + ExecutionObject execObject = transformerService.transform(flow.getFlow()); + UUID id = UUID.randomUUID(); + execObject.setId(id.toString()); + executions.put(execObject.getId(), execObject); + return execObject; + } + + public ExecutionObject getExecution(String id) { + return executions.get(id); + } + + public List getAllExecutions() { + return executions.values().stream().toList(); + } + + public void removeExecution(String id) { + executions.remove(id); + } + + public void prepareInput(String id, String key, Object input) { + ExecutionObject eo = getExecution(id); + if (eo == null) throw new IllegalArgumentException("Execution with id "+id+" not found"); + eo.getContext().setStatus(Status.INITIALIZING); + eo.getStartStep().inputReady(key, input); + if (eo.getStartStep().areAllInputsReady()) { + eo.getContext().setStatus(Status.READY); + } + } + + public void startExecution(String id) { + ExecutionObject eo = getExecution(id); + if (eo == null) throw new IllegalArgumentException("Execution with id "+id+" not found"); + if (eo.getContext().getStatus() != Status.READY) { + throw new IllegalStateException("Execution with id "+id+" is not ready"); + } + eo.getStartStep().startExecution(); + eo.getContext().setStatus(Status.RUNNING); + } + + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java b/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java index efb9b72..7b367bc 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java @@ -4,8 +4,6 @@ import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Objects; -import java.util.stream.Stream; import org.slf4j.Logger; import org.springframework.beans.factory.annotation.Autowired; @@ -13,6 +11,7 @@ import org.springframework.stereotype.Service; 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.ExecutionStep; import it.cnr.isti.workflow.manager.model.types.NodeDefinition; import it.cnr.isti.workflow.manager.model.ui.Connection; @@ -34,13 +33,16 @@ public class TransformerService { private ExecutionObject execObject; - ExecutionObject transform(Flow flow) { + public ExecutionObject transform(Flow flow) { execObject = new ExecutionObject(); - execObject.setStartStep(createStep(flow, null, flow.getStartNode(), null)); + ExecutionContext context = new ExecutionContext(); + execObject.setContext(context); + execObject.setStartStep(createStep(flow, null, context, flow.getStartNode(), null)); + execObject.getStartStep().setStartAutomatically(false); return execObject; } - private ExecutionStep createStep(Flow flow, ExecutionStep parent, String nodeId, Map mappingNodeOutputNameToChildInputName ) { + private ExecutionStep createStep(Flow flow, ExecutionStep parent, ExecutionContext context, String nodeId, Map mappingNodeOutputNameToChildInputName ) { //TODO: cosa succede se ho input da più nodi in un solo nodo ? devo tenere traccia dei nodi già visitati e non ricreare lo step ma aggiungerlo solamente ai previous steps @@ -58,7 +60,7 @@ public class TransformerService { throw new IllegalArgumentException("Executor not found for node type: " + node.getType()); ExecutionStep step = ExecutionStep.builder().id(node.getKey()).executor(executors.get(nodeDefinition.getExecutor())) - .runtimeParameters(node.getParameters()).nodeDefinition(nodeDefinition).build(); + .runtimeParameters(node.getParameters()).nodeDefinition(nodeDefinition).context(context).build(); if (parent == null){ // se è il primo step devo settare il mapping degli input Map mappingInputName = new HashMap<>(); @@ -90,7 +92,7 @@ public class TransformerService { } } //creo ricorsivamente i prossimi step - step.getNextSteps().addAll(childNodes.stream().map(n -> createStep(flow, step, n.getKey(), currentMappingNodeOutputNameToChildInputName)).toList()); + step.getNextSteps().addAll(childNodes.stream().map(n -> createStep(flow, step, context, n.getKey(), currentMappingNodeOutputNameToChildInputName)).toList()); } else execObject.setEndStep(step); 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 e8e9c36..c0a4985 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java +++ b/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java @@ -6,7 +6,6 @@ import org.springframework.boot.CommandLineRunner; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration; import org.springframework.context.annotation.Bean; import it.cnr.isti.workflow.manager.model.types.Key; @@ -29,33 +28,17 @@ public class WorkflowManagerApplication { @ConditionalOnProperty(prefix = "app", name = "db.init.enabled", havingValue = "true") CommandLineRunner init(NodeDefinitionRepository repository) { return args -> { - NodeDefinition defectDetection = NodeDefinition.builder().type("DefectDetection").name("defect detection") - .executor("GENERIC-AI").category("DefectDetection").input("file","CSV").output("file","CSV") - .build(); - + JSONObject rPars = new JSONObject().put("options", - List.of(new JSONObject().put("key", "1").put("value", "ChatGpt"), - new JSONObject().put("key", "2").put("value", "LLama"))); + List.of("google-gemini", "OLLAMA")); ParameterDefinition rgParams = ParameterDefinition.builder().name("LLM").type(ParameterType.Select) .label("Select the LLM to use").description("The LLM to use for generating requirements") .validation(Validation.builder().name("required").validator("required").build()) .specificAttributes(rPars).build(); - NodeDefinition requirementGeneration = NodeDefinition.builder().type("RequirementGeneration") - .executor("GENERIC-AI").name("requirement generation").category("RequirementGeneration") - .input("file","CSV").output("file","CSV") - .runtimeParameter("LLM", rgParams).build(); - - /*ParameterDefinition phraseParam = ParameterDefinition.builder().name("phrase").type(ParameterType.Input) - .label("Phrase").description("The prhase to translate") - .validation(Validation.builder().name("required").validator("required").build()) - .build(); */ - JSONObject languageOptions = new JSONObject().put("options", - List.of(new JSONObject().put("key", "1").put("value", "French"), - new JSONObject().put("key", "2").put("value", "English"))); - + List.of("FR", "EN", "IT", "ES", "DE")); ParameterDefinition langParam = ParameterDefinition.builder().name("language").type(ParameterType.Select) .label("Language").description("The language to translate to") .validation(Validation.builder().name("required").validator("required").build()) @@ -69,14 +52,11 @@ public class WorkflowManagerApplication { .executor("GENERIC-AI").name("translator").category("TRANSLATOR") .input("phrase","Text").output("translated","Text") .runtimeParameter("LLM", rgParams).runtimeParameter("language", langParam) - .executorToOutputTraslationMapping("response", "translated") + .executorToOutputTranslationMapping("response", "translated") .inputTranslator("prompt", inputTranslator).build(); - repository.save(defectDetection); - repository.save(requirementGeneration); repository.save(translator); - }; } 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 new file mode 100644 index 0000000..a7efce3 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ExecutionController.java @@ -0,0 +1,92 @@ +package it.cnr.isti.workflow.manager.controllers; + +import java.io.File; +import java.io.IOException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.web.server.WebServerException; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.CrossOrigin; +import org.springframework.web.bind.annotation.ExceptionHandler; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestParam; +import org.springframework.web.bind.annotation.ResponseBody; +import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.multipart.MultipartFile; + +import it.cnr.isti.workflow.manager.ExecutionService; +import it.cnr.isti.workflow.manager.TransformerService; +import it.cnr.isti.workflow.manager.executors.ExecutionObject; +import it.cnr.isti.workflow.manager.model.ui.FlowEntity; +import it.cnr.isti.workflow.manager.repositories.FlowRepository; +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.PathVariable; + +@RestController +@CrossOrigin(origins = "http://localhost:4200") +@RequestMapping("/executions") +public class ExecutionController { + + @Autowired + FlowRepository flowRepository; + + @Autowired + ExecutionService executionService; + + @PostMapping() + public String create(@RequestBody String flowId) { + FlowEntity flow = flowRepository.findById(flowId) + .orElseThrow(() -> new IllegalArgumentException("Flow with id " + flowId + " not found")); + ExecutionObject eo = executionService.createExecution(flow); + return eo.getId(); + } + + @GetMapping() + public List getAll() { + return executionService.getAllExecutions(); + } + + @GetMapping(path = "{id}") + public ExecutionObject get(@PathVariable String id) { + return executionService.getExecution(id); + } + + @PutMapping(path = "{id}/start") + public void start(@PathVariable String id) { + executionService.startExecution(id); + } + + @PutMapping(path = "{id}/input/{inputName}") + public void prepareStringInputs(@RequestBody String input, @PathVariable String inputName, + @PathVariable String id) { + executionService.prepareInput(id, inputName, input); + } + + + @PutMapping(path = "{id}/input/{inputName}", consumes = "multipart/form-data") + @ResponseBody + public ResponseEntity prepareFileInputs(@RequestParam("file") MultipartFile file, @PathVariable String inputName, + @PathVariable String id) { + try { + File myFile = File.createTempFile(inputName, file.getOriginalFilename()); + file.transferTo(myFile); + executionService.prepareInput(id, inputName, myFile); + } catch (IOException e) { + throw new WebServerException("Error while creating file", e); + } + return ResponseEntity.ok().build(); + + } + + @ExceptionHandler(IllegalArgumentException.class) + public ResponseEntity executionNotFound(IllegalArgumentException exc) { + return ResponseEntity.notFound().build(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java index b1af4f0..19fd71b 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/NodeDefinitionController.java @@ -1,6 +1,7 @@ package it.cnr.isti.workflow.manager.controllers; import java.util.List; +import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.CrossOrigin; @@ -10,9 +11,12 @@ import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; +import it.cnr.isti.workflow.manager.executors.Executor; +import it.cnr.isti.workflow.manager.model.ExecutorDescriptor; import it.cnr.isti.workflow.manager.model.types.NodeDefinition; import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository; + @RestController @CrossOrigin(origins = "http://localhost:4200") @RequestMapping("/types") @@ -22,6 +26,9 @@ public class NodeDefinitionController { @Autowired private NodeDefinitionRepository repository; + @Autowired + private Map executors; + @GetMapping("/nodes") public List getNodes() { return (List) repository.findAll(); @@ -31,6 +38,11 @@ public class NodeDefinitionController { void addNode(@RequestBody NodeDefinition node) { repository.save(node); } - + + @GetMapping("/executors") + public List getAvailableExecutors() { + return executors.entrySet().stream() + .map((entry) -> new ExecutorDescriptor(entry.getKey(), entry.getValue())).toList(); + } } 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 0047a28..5e9e4bd 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 @@ -1,13 +1,54 @@ package it.cnr.isti.workflow.manager.executors; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import com.fasterxml.jackson.annotation.JsonIgnore; + +import it.cnr.isti.workflow.manager.model.ExecutionContext; import it.cnr.isti.workflow.manager.model.ExecutionStep; +import it.cnr.isti.workflow.manager.model.InputReadyListener; +import it.cnr.isti.workflow.manager.model.ExecutionContext.Status; import lombok.Data; import lombok.NoArgsConstructor; @Data @NoArgsConstructor -public class ExecutionObject { +public class ExecutionObject implements InputReadyListener { + String id; + String flowName; + + + ExecutionContext context; + + @JsonIgnore ExecutionStep startStep; + + @JsonIgnore ExecutionStep endStep; + + public void setEndStep(ExecutionStep endStep) { + endStep.setNextSteps(List.of(this)); + this.endStep = endStep; + + } + + @JsonIgnore + Map inputs = new HashMap<>(); + + @Override + public void inputReady(String key, Object value) { + if (this.context.getResult() == null) { + this.context.setResult(new HashMap<>()); + } + this.context.getResult().put(key, value); + if (this.endStep.getNodeDefinition().getOutputs().size() == this.context.getResult().size() ){ + //significa che l'esecuzione è finita + this.context.setStatus(Status.SUCCESS); + } + + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/Executor.java b/src/main/java/it/cnr/isti/workflow/manager/executors/Executor.java index 26c0e15..f9f9146 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/Executor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/Executor.java @@ -1,12 +1,21 @@ package it.cnr.isti.workflow.manager.executors; +import java.util.List; import java.util.Map; +import it.cnr.isti.workflow.manager.model.types.ParameterDefinition; + public interface Executor { Map execute(Map userParameters, Map preExecutionParameters, Map inputsFromParent); Map> getRequiredPreExecutionParameters(Map userParameters); + List declaredOutputNames(); + + List declaredInputNames(); + + List getMandatoryParameters(); + String getName(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GenericAIExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GenericAIExecutor.java index ddea185..81742a0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GenericAIExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GenericAIExecutor.java @@ -1,12 +1,16 @@ package it.cnr.isti.workflow.manager.executors.ai; +import java.util.List; import java.util.Map; import java.util.Objects; - +import org.json.JSONObject; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import it.cnr.isti.workflow.manager.executors.Executor; +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 lombok.ToString; @Service("GENERIC-AI") @@ -49,30 +53,31 @@ public class GenericAIExecutor implements Executor { return aiModels.get(userParameters.get("LLM")).getRequiredPreExecutionParameters(); } + @Override + public List declaredOutputNames() { + return List.of("response"); + } + + @Override + public List declaredInputNames() { + return List.of("prompt"); + } + + @Override + public List getMandatoryParameters() { + + JSONObject rPars = new JSONObject().put("options", + aiModels.keySet()); + + + ParameterDefinition rgParams = ParameterDefinition.builder().name("LLM").type(ParameterType.Select) + .label("Select the LLM to use").description("The LLM to use for generating requirements") + .validation(Validation.builder().name("required").validator("required").build()) + .specificAttributes(rPars).build(); + return List.of(rgParams); + } + + - - /* - @Override - public String onExecute(Map parameters, Map executorSpecificParameters, - String payload) { - return "AI Executor: " + payload; - } - - @Override - public boolean areParametersValid(Map parameters) { - return parameters != null && parameters.containsKey("model") && aiModels.containsKey(parameters.get("model")); - } - - @Override - public String mapsGenericInputsToPayload(Map parameters) { - return (String) parameters.get("prompt"); - } - - - @Override - public Map> getExecutorSpecificParametersName() { - return Map.of("model", String.class); - } */ - } 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 new file mode 100644 index 0000000..5626968 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionContext.java @@ -0,0 +1,33 @@ +package it.cnr.isti.workflow.manager.model; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +public class ExecutionContext { + + public enum Status { + CREATED, + INITIALIZING, + READY, + RUNNING, + SUCCESS, + ERROR + } + + Map result = null; + + List stepsUnderExecution = new ArrayList<>(); + List errors = new ArrayList<>();; + List warnings = new ArrayList<>();; + + Status status = Status.CREATED; + +} + + diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionStep.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionStep.java index 90a5d7a..4e3d393 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionStep.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutionStep.java @@ -29,6 +29,12 @@ public class ExecutionStep implements InputReadyListener { String id; Map runtimeParameters; Map fixedParameters; + + @NonNull + ExecutionContext context; + + boolean startAutomatically = true; + Executor executor; Map parentOutputsToInputMapping; @@ -42,7 +48,8 @@ public class ExecutionStep implements InputReadyListener { NodeDefinition nodeDefinition; @Builder - public ExecutionStep(String id, Map runtimeParameters, Map fixedParameters, + public ExecutionStep(ExecutionContext context, String id, Map runtimeParameters, + Map fixedParameters, NodeDefinition nodeDefinition, Executor executor) { this.id = id; @@ -50,57 +57,86 @@ public class ExecutionStep implements InputReadyListener { this.fixedParameters = fixedParameters; this.executor = executor; this.nodeDefinition = nodeDefinition; + this.context = context; } @Override public void inputReady(String inputKey, Object inputValue) { log.debug("Input ready for step {} with key {} and value {}", id, inputKey, inputValue); - parentOutputsToInputMapping.forEach((k,v) -> log.debug("mapping {} -> {}", k, v)); + parentOutputsToInputMapping.forEach((k, v) -> log.debug("mapping {} -> {}", k, v)); if (parentOutputsToInputMapping.containsKey(inputKey)) { inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue); - if (areAllInputsReady()) { - log.debug("all inputs ready for step {}", id); - 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(); - - for (Key usedKey : entry.getValue().getUsedKeys()) { - String valueToReplace = ""; - switch (usedKey.getType()) { - case Key.Type.RUNTIME: - valueToReplace = (String) this.runtimeParameters.get(usedKey.getKey()); - break; - default: - valueToReplace = (String) inputs.get(usedKey.getKey()); - keys.remove(usedKey.getKey()); - break; - } - translation = translation.replace(String.format("${{%s}}", usedKey.getKey()), valueToReplace); - log.debug("replacing used key {} value {}", usedKey.getKey(), valueToReplace); - } - realInputs.put(entry.getKey(), translation); - } - keys.forEach(k -> realInputs.put(k, inputs.get(k))); - } else - realInputs.putAll(inputs); - - log.debug("input translation for step {} is {}", id, realInputs); - Map returned = executor.execute(this.runtimeParameters, null, realInputs); - for (InputReadyListener nextStep : nextSteps) - for (Map.Entry entry : returned.entrySet()) - nextStep.inputReady(!nodeDefinition.getExecutorToOutputTraslationMappings().isEmpty() - ? nodeDefinition.getExecutorToOutputTraslationMappings().get(entry.getKey()) - : entry.getKey(), entry.getValue()); - } else { - log.debug("inputs not ready for step {}", id); - } + if (startAutomatically) startExecution(); } }; - boolean areAllInputsReady() { + public boolean areAllInputsReady() { return inputs.size() == parentOutputsToInputMapping.size(); } + + public void startExecution() { + if (areAllInputsReady()) + new Thread(() -> { + this.context.getStepsUnderExecution().add(this.id); + try { + Map realInputs = prepareInputs(); + log.debug("input translation for step {} is {}", id, realInputs); + Map returned = executor.execute(this.runtimeParameters, null, realInputs); + for (InputReadyListener nextStep : nextSteps) + for (Map.Entry entry : returned.entrySet()) + nextStep.inputReady(!nodeDefinition.getExecutorToOutputTranslationMappings().isEmpty() + ? nodeDefinition.getExecutorToOutputTranslationMappings().get(entry.getKey()) + : entry.getKey(), entry.getValue()); + } catch (Throwable e) { + log.error("error during execution of step {}: {}", id, e.getMessage()); + this.context.getErrors().add("[" + this.id + "] e.getMessage()"); + } + this.context.getStepsUnderExecution().remove(this.id); + + }).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(); + + for (Key usedKey : entry.getValue().getUsedKeys()) { + String valueToReplace = ""; + switch (usedKey.getType()) { + case Key.Type.RUNTIME: + valueToReplace = (String) this.runtimeParameters.get(usedKey.getKey()); + break; + default: + valueToReplace = (String) inputs.get(usedKey.getKey()); + keys.remove(usedKey.getKey()); + break; + } + translation = translation.replace(String.format("${{%s}}", usedKey.getKey()), valueToReplace); + log.debug("replacing used key {} value {}", usedKey.getKey(), 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/ExecutorDescriptor.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java new file mode 100644 index 0000000..95c437b --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java @@ -0,0 +1,30 @@ +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.Data; + +@Data +public class ExecutorDescriptor { + + + public ExecutorDescriptor(String id, Executor executor){ + this.identifier = id; + this.name = executor.getName(); + this.inputNames = executor.declaredInputNames(); + this.outputNames = executor.declaredOutputNames(); + this.mandatoryParameters = executor.getMandatoryParameters(); + } + + String identifier; + String name; + String description; + List inputNames; + List outputNames; + List mandatoryParameters; + + //List optionalParameters; + +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/NodeDefinition.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/NodeDefinition.java index 5eaf888..2e0523a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/model/types/NodeDefinition.java +++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/NodeDefinition.java @@ -1,13 +1,10 @@ package it.cnr.isti.workflow.manager.model.types; -import java.util.List; import java.util.Map; - import it.cnr.isti.workflow.manager.repositories.InputTranslatorConverter; import it.cnr.isti.workflow.manager.repositories.MapsConverter; import jakarta.persistence.Column; import jakarta.persistence.Convert; -import jakarta.persistence.Embedded; import jakarta.persistence.Entity; import jakarta.persistence.Id; import lombok.AllArgsConstructor; @@ -61,8 +58,15 @@ public class NodeDefinition { @Singular @Convert(converter = MapsConverter.class) @Column(columnDefinition = "TEXT") - private Map executorToOutputTraslationMappings; + private Map executorToOutputTranslationMappings; + //key is the executor input name, value is the input name + @Singular + @Convert(converter = MapsConverter.class) + @Column(columnDefinition = "TEXT") + private Map inputToExecutorTranslationMappings; + + //key is the executor output name, value is the output name @Singular @Convert(converter = InputTranslatorConverter.class) @Column(columnDefinition = "TEXT") diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index c01c016..545956b 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -1,9 +1,9 @@ spring.application.name=workflow-manager #postgresql details -spring.datasource.url=jdbc:postgresql://localhost:5432/mydb -spring.datasource.username=lucio -spring.datasource.password=password +spring.datasource.url=${DB_URL} +spring.datasource.username=${DB_USER} +spring.datasource.password=${DB_PASSWORD} spring.datasource.driver-class-name=org.postgresql.Driver spring.jpa.properties.hibernate.dialect=org.hibernate.dialect.PostgreSQLDialect