execution controller added
This commit is contained in:
parent
d3e7922e38
commit
fcd83c04d0
|
|
@ -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"]
|
||||
13
pom.xml
13
pom.xml
|
|
@ -12,7 +12,7 @@
|
|||
<groupId>it.cnr.isti</groupId>
|
||||
<artifactId>workflow-manager</artifactId>
|
||||
<version>0.0.1-SNAPSHOT</version>
|
||||
<packaging>war</packaging>
|
||||
<packaging>jar</packaging>
|
||||
<name>workflow-manager</name>
|
||||
<description>Workflow server project</description>
|
||||
<url />
|
||||
|
|
@ -49,17 +49,17 @@
|
|||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-webflux</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springdoc</groupId>
|
||||
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
|
||||
<version>2.0.2</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.projectlombok</groupId>
|
||||
<artifactId>lombok</artifactId>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-tomcat</artifactId>
|
||||
<scope>provided</scope>
|
||||
</dependency>
|
||||
<!-- https://mvnrepository.com/artifact/org.json/json -->
|
||||
<dependency>
|
||||
<groupId>org.json</groupId>
|
||||
|
|
@ -75,6 +75,7 @@
|
|||
</dependencies>
|
||||
|
||||
<build>
|
||||
<finalName>workflow-manager</finalName>
|
||||
<plugins>
|
||||
<plugin>
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
|
|
|
|||
|
|
@ -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<String,ExecutionObject> 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<ExecutionObject> 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);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
|
@ -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<String, String> mappingNodeOutputNameToChildInputName ) {
|
||||
private ExecutionStep createStep(Flow flow, ExecutionStep parent, ExecutionContext context, String nodeId, Map<String, String> 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<String, String> 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);
|
||||
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
||||
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<ExecutionObject> 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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String,Executor> executors;
|
||||
|
||||
@GetMapping("/nodes")
|
||||
public List<NodeDefinition> getNodes() {
|
||||
return (List<NodeDefinition>) repository.findAll();
|
||||
|
|
@ -31,6 +38,11 @@ public class NodeDefinitionController {
|
|||
void addNode(@RequestBody NodeDefinition node) {
|
||||
repository.save(node);
|
||||
}
|
||||
|
||||
|
||||
@GetMapping("/executors")
|
||||
public List<ExecutorDescriptor> getAvailableExecutors() {
|
||||
return executors.entrySet().stream()
|
||||
.map((entry) -> new ExecutorDescriptor(entry.getKey(), entry.getValue())).toList();
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String,Object> execute(Map<String, Object> userParameters, Map<String,Object> preExecutionParameters, Map<String, Object> inputsFromParent);
|
||||
|
||||
Map<String, Class<?>> getRequiredPreExecutionParameters(Map<String, Object> userParameters);
|
||||
|
||||
List<String> declaredOutputNames();
|
||||
|
||||
List<String> declaredInputNames();
|
||||
|
||||
List<ParameterDefinition> getMandatoryParameters();
|
||||
|
||||
String getName();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String> declaredOutputNames() {
|
||||
return List.of("response");
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> declaredInputNames() {
|
||||
return List.of("prompt");
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<ParameterDefinition> 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<String, Object> parameters, Map<String, Object> executorSpecificParameters,
|
||||
String payload) {
|
||||
return "AI Executor: " + payload;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean areParametersValid(Map<String, Object> parameters) {
|
||||
return parameters != null && parameters.containsKey("model") && aiModels.containsKey(parameters.get("model"));
|
||||
}
|
||||
|
||||
@Override
|
||||
public String mapsGenericInputsToPayload(Map<String, Object> parameters) {
|
||||
return (String) parameters.get("prompt");
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public Map<String, Class<?>> getExecutorSpecificParametersName() {
|
||||
return Map.of("model", String.class);
|
||||
} */
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object> result = null;
|
||||
|
||||
List<String> stepsUnderExecution = new ArrayList<>();
|
||||
List<String> errors = new ArrayList<>();;
|
||||
List<String> warnings = new ArrayList<>();;
|
||||
|
||||
Status status = Status.CREATED;
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -29,6 +29,12 @@ public class ExecutionStep implements InputReadyListener {
|
|||
String id;
|
||||
Map<String, Object> runtimeParameters;
|
||||
Map<String, Object> fixedParameters;
|
||||
|
||||
@NonNull
|
||||
ExecutionContext context;
|
||||
|
||||
boolean startAutomatically = true;
|
||||
|
||||
Executor executor;
|
||||
Map<String, String> parentOutputsToInputMapping;
|
||||
|
||||
|
|
@ -42,7 +48,8 @@ public class ExecutionStep implements InputReadyListener {
|
|||
NodeDefinition nodeDefinition;
|
||||
|
||||
@Builder
|
||||
public ExecutionStep(String id, Map<String, Object> runtimeParameters, Map<String, Object> fixedParameters,
|
||||
public ExecutionStep(ExecutionContext context, String id, Map<String, Object> runtimeParameters,
|
||||
Map<String, Object> 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<String, Object> realInputs = new HashMap<>();
|
||||
Set<String> keys = new HashSet<>(inputs.keySet());
|
||||
if (nodeDefinition.getInputTranslators() != null) {
|
||||
log.debug("translators {}", nodeDefinition.getInputTranslators().size());
|
||||
for (Map.Entry<String, Translator> 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<String, Object> returned = executor.execute(this.runtimeParameters, null, realInputs);
|
||||
for (InputReadyListener nextStep : nextSteps)
|
||||
for (Map.Entry<String, Object> 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<String, Object> realInputs = prepareInputs();
|
||||
log.debug("input translation for step {} is {}", id, realInputs);
|
||||
Map<String, Object> returned = executor.execute(this.runtimeParameters, null, realInputs);
|
||||
for (InputReadyListener nextStep : nextSteps)
|
||||
for (Map.Entry<String, Object> 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<String, Object> prepareInputs() {
|
||||
Map<String, Object> realInputs = new HashMap<>();
|
||||
Set<String> keys = new HashSet<>(inputs.keySet());
|
||||
if (nodeDefinition.getInputTranslators() != null) {
|
||||
log.debug("translators {}", nodeDefinition.getInputTranslators().size());
|
||||
for (Map.Entry<String, Translator> 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<String, String> 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;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String> inputNames;
|
||||
List<String> outputNames;
|
||||
List<ParameterDefinition> mandatoryParameters;
|
||||
|
||||
//List<ParameterDescriptor> optionalParameters;
|
||||
|
||||
}
|
||||
|
|
@ -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<String, String> executorToOutputTraslationMappings;
|
||||
private Map<String, String> executorToOutputTranslationMappings;
|
||||
|
||||
//key is the executor input name, value is the input name
|
||||
@Singular
|
||||
@Convert(converter = MapsConverter.class)
|
||||
@Column(columnDefinition = "TEXT")
|
||||
private Map<String, String> inputToExecutorTranslationMappings;
|
||||
|
||||
//key is the executor output name, value is the output name
|
||||
@Singular
|
||||
@Convert(converter = InputTranslatorConverter.class)
|
||||
@Column(columnDefinition = "TEXT")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue