diff --git a/pom.xml b/pom.xml
index 88e796f..5f337c2 100644
--- a/pom.xml
+++ b/pom.xml
@@ -85,6 +85,11 @@
spring-boot-starter-test
test
+
+ com.h2database
+ h2
+ test
+
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 16744c3..d678784 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/TransformerService.java
@@ -31,27 +31,53 @@ public class TransformerService {
@Autowired
NodeDefinitionRepository nodeDefinitionRepository;
- private ExecutionObject execObject;
-
public ExecutionObject transform(Flow flow) {
- execObject = new ExecutionObject(flow);
- ExecutionStep startStep = createStep(flow, null, execObject.getContext(), flow.getStartNode(), null);
- execObject.setStartStep(startStep);
+
+ Map steps = new HashMap<>();
+ Map inputs = new HashMap<>();
+ Map outputs = new HashMap<>();
+
+ ExecutionObject execObject = new ExecutionObject(flow);
+
+ for (Node node : flow.getNodes()) {
+ ExecutionStep step = getStep(node, execObject.getContext());
+ steps.put(node.getKey(), step);
+
+ if (node.getKey().equals(flow.getStartNode())) {
+ Map mappingInputName = new HashMap<>();
+ step.getNodeDefinition().getInputs().forEach((k, v) -> mappingInputName.put(k, k));
+ log.debug("setting {} to step {}", mappingInputName, step.getId());
+ step.setParentOutputsToInputMapping(mappingInputName);
+ execObject.setStartStep(step);
+ node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName())));
+ } else if (node.getKey().equals(flow.getEndNode())) {
+ execObject.setEndStep(step);
+ node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName())));
+ } else {
+ 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()) {
+ String from = connection.getFrom();
+ String to = connection.getTo();
+ if (!inputs.containsKey(to))
+ throw new RuntimeException("Input " + to + " not found");
+ if (!outputs.containsKey(from))
+ throw new RuntimeException("Output " + from + " not found");
+ IDPair inputPair = inputs.get(to);
+ IDPair outputPair = outputs.get(from);
+ ExecutionStep outputNode = steps.get(outputPair.nodeId);
+ ExecutionStep inputNode = steps.get(inputPair.nodeId);
+ outputNode.getNextSteps().add(inputNode);
+ inputNode.getParentOutputsToInputMapping().put(outputPair.name, inputPair.name);
+ }
+
return execObject;
}
- 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
-
- Node node = flow.getNodes().stream().filter(n -> n.getKey().equals(nodeId)).findFirst()
- .orElseThrow(() -> new RuntimeException("node " + nodeId + " not found"));
-
- // prendo tutte le gli id in uscita dal nodo
- List outputs = node.getOutputs();
-
+ private ExecutionStep getStep(Node node, ExecutionContext context) {
// prende l'executor per il tipo di nodo
NodeDefinition nodeDefinition = nodeDefinitionRepository.findById(node.getType())
.orElseThrow(() -> new RuntimeException("Node definition not found for node type: " + node.getType()));
@@ -63,48 +89,16 @@ public class TransformerService {
.executor(executors.get(nodeDefinition.getExecutor()))
.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<>();
- nodeDefinition.getInputs().forEach((k, v) -> mappingInputName.put(k, k));
- log.debug("setting {} to step {}", mappingInputName, step.getId());
- step.setParentOutputsToInputMapping(mappingInputName);
- } else {
- step.setParentOutputsToInputMapping(mappingNodeOutputNameToChildInputName);
- }
-
- if (!nodeId.equals(flow.getEndNode())) {
-
- // mantiene il mapping tra nome output del nodo inviante e nome input del nodo
- // ricevente
- Map currentMappingNodeOutputNameToChildInputName = new HashMap<>();
-
- List childNodes = new ArrayList<>();
-
- for (IOModel output : outputs) {
- for (Connection connection : flow.getConnections()) {
- if (!connection.getFrom().equals(output.getKey()))
- continue;
- String inputId = connection.getTo();
- Node nodeFound = flow.getNodes().stream()
- .filter(n -> n.getInputs().stream().anyMatch(i -> i.getKey().equals(inputId))).findFirst()
- .orElseThrow(() -> new RuntimeException("node " + inputId + " not found"));
-
- // prendo il nome dell'input e lo mappo con l'output
- IOModel input = nodeFound.getInputs().stream().filter(i -> i.getKey().equals(inputId)).findFirst()
- .orElseThrow(() -> new RuntimeException("input " + inputId + " not found"));
- currentMappingNodeOutputNameToChildInputName.put(output.getName(), input.getName());
- childNodes.add(nodeFound);
- }
- }
- // creo ricorsivamente i prossimi step
- step.getNextSteps().addAll(childNodes.stream()
- .map(n -> createStep(flow, step, context, n.getKey(), currentMappingNodeOutputNameToChildInputName))
- .toList());
- } else
- execObject.setEndStep(step);
-
return step;
-
}
+ private static class IDPair {
+ public String nodeId;
+ public String name;
+
+ public IDPair(String nodeId, String name) {
+ this.nodeId = nodeId;
+ this.name = name;
+ }
+ }
}
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 9833f98..8907b7f 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/WorkflowManagerApplication.java
@@ -8,6 +8,9 @@ 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.util.ResourceUtils;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.executors.ai.AIModel;
import it.cnr.isti.workflow.manager.model.types.IOType;
@@ -16,7 +19,9 @@ import it.cnr.isti.workflow.manager.model.types.NodeDefinition;
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.Translator;
-import it.cnr.isti.workflow.manager.model.types.Validation;
+import it.cnr.isti.workflow.manager.model.ui.Flow;
+import it.cnr.isti.workflow.manager.model.ui.FlowEntity;
+import it.cnr.isti.workflow.manager.repositories.FlowRepository;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
@SpringBootApplication
@@ -31,7 +36,7 @@ public class WorkflowManagerApplication {
@Bean
@ConditionalOnProperty(prefix = "app", name = "db.init.enabled", havingValue = "true")
- CommandLineRunner init(NodeDefinitionRepository repository) {
+ CommandLineRunner init(NodeDefinitionRepository repository, FlowRepository flowRepository) {
return args -> {
JSONObject rPars = new JSONObject().put("options",
@@ -39,18 +44,16 @@ public class WorkflowManagerApplication {
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();
+ .required(true).specificAttributes(rPars).build();
JSONObject languageOptions = new JSONObject().put("options",
List.of("Franch", "English", "Italian", "spanish", "German", "Chinese", "Japanese"));
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())
+ .label("Language").description("The language to translate to").required(true)
.specificAttributes(languageOptions).build();
- Translator inputTranslator = Translator.builder().usedKey(Key.inputKey("phrase")).usedKey(Key.runtimeKey("language"))
+ Translator inputTranslator = Translator.builder()
.translation("translate the following phrase ${{phrase}} to ${{language}} returning only the translated").build();
NodeDefinition translator = NodeDefinition.builder().type("Translator")
@@ -63,9 +66,39 @@ public class WorkflowManagerApplication {
repository.save(translator);
+ loadNodeDefinitionAndSave(repository);
+
+ loadFlowsAndSave(flowRepository);
};
}
+ private void loadNodeDefinitionAndSave(NodeDefinitionRepository repository) {
+ ObjectMapper objectMapper = new ObjectMapper();
+
+ try{
+
+ List nodes = objectMapper.readValue(ResourceUtils.getFile("classpath:saved-node-type.json"), objectMapper.getTypeFactory().constructCollectionType(List.class,NodeDefinition.class));
+ for (NodeDefinition node: nodes)
+ repository.save(node);
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ }
+ private void loadFlowsAndSave(FlowRepository repository) {
+ ObjectMapper objectMapper = new ObjectMapper();
+
+ try{
+
+ List flows = objectMapper.readValue(ResourceUtils.getFile("classpath:flows.json"), objectMapper.getTypeFactory().constructCollectionType(List.class,FlowEntity.class));
+ for (FlowEntity flow: flows){
+ FlowEntity savedFlow = repository.save(flow);
+ savedFlow.getFlow().setId(savedFlow.getId());
+ repository.save(savedFlow);
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ }
}
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 9341962..97fbfba 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
@@ -6,7 +6,6 @@ import java.util.List;
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;
@@ -27,6 +26,8 @@ import org.springframework.web.bind.annotation.PathVariable;
@RequestMapping("/executions")
public class ExecutionController {
+ private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ExecutionController.class);
+
@Autowired
FlowRepository flowRepository;
@@ -35,10 +36,16 @@ public class ExecutionController {
@PostMapping()
public ExecutionObject create(@RequestBody String flowId) {
+ logger.info("Creating execution for flow {}", flowId);
FlowEntity flow = flowRepository.findById(flowId)
.orElseThrow(() -> new IllegalArgumentException("Flow with id " + flowId + " not found"));
- ExecutionObject eo = executionService.createExecution(flow);
- return eo;
+ try{
+ ExecutionObject eo = executionService.createExecution(flow);
+ return eo;
+ } catch (Throwable e) {
+ logger.error("Error creating execution for flow {}", flowId, e);
+ throw new WebServerException("Error while creating execution", e);
+ }
}
@GetMapping()
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 006a0de..7774995 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,11 +1,9 @@
package it.cnr.isti.workflow.manager.controllers;
import java.util.List;
-import java.util.Map;
import org.slf4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.web.bind.annotation.CrossOrigin;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
@@ -27,7 +25,7 @@ public class NodeDefinitionController {
private NodeDefinitionRepository repository;
@Autowired
- private Map executors;
+ private List executors;
@GetMapping("/nodes")
public List getNodes() {
@@ -51,8 +49,7 @@ public class NodeDefinitionController {
@GetMapping("/executors")
public List getAvailableExecutors() {
- return executors.entrySet().stream()
- .map((entry) -> new ExecutorDescriptor(entry.getKey(), entry.getValue())).toList();
+ return executors.stream().map(entry -> entry.getDescriptor(null)).toList();
}
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/exceptions/ExecutionException.java b/src/main/java/it/cnr/isti/workflow/manager/exceptions/ExecutionException.java
new file mode 100644
index 0000000..8cccd6b
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/exceptions/ExecutionException.java
@@ -0,0 +1,24 @@
+package it.cnr.isti.workflow.manager.exceptions;
+
+public class ExecutionException extends Exception {
+
+ private static final long serialVersionUID = 1L;
+
+ public ExecutionException(String message) {
+ super(message);
+ }
+
+ public ExecutionException(String message, Throwable cause) {
+ super(message, cause);
+ }
+
+ public ExecutionException(Throwable cause) {
+ super(cause);
+ }
+
+ public ExecutionException(String message, Throwable cause, boolean enableSuppression,
+ boolean writableStackTrace) {
+ super(message, cause, enableSuppression, writableStackTrace);
+ }
+
+}
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 f9f9146..d83ceb2 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,21 +1,15 @@
package it.cnr.isti.workflow.manager.executors;
-import java.util.List;
import java.util.Map;
-import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
+import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
+import it.cnr.isti.workflow.manager.model.ExecutorDescriptor;
public interface Executor {
- Map execute(Map userParameters, Map preExecutionParameters, Map inputsFromParent);
+ Map execute(Map userParameters, Map preExecutionParameters, Map inputsFromParent) throws ExecutionException;
- Map> getRequiredPreExecutionParameters(Map userParameters);
-
- List declaredOutputNames();
-
- List declaredInputNames();
-
- List getMandatoryParameters();
-
+ ExecutorDescriptor getDescriptor(Map executorDescriptorParameters);
+
String getName();
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java
index 3bb596a..d037d75 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executors/NoOpExecutor.java
@@ -1,21 +1,34 @@
package it.cnr.isti.workflow.manager.executors;
import org.springframework.stereotype.Service;
+
+import java.lang.reflect.Parameter;
import java.util.List;
import java.util.Map;
+
+import it.cnr.isti.workflow.manager.model.ExecutorDescriptor;
+import it.cnr.isti.workflow.manager.model.Validations;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
+import it.cnr.isti.workflow.manager.model.types.ParameterType;
import lombok.ToString;
@Service("NoOp")
@ToString(of = {"name"})
public class NoOpExecutor implements Executor {
- private static final String name = "NoOp";
+
+ private static final String name = "No Operation Executor";
@Override
public Map execute(Map userParameters, Map preExecutionParameters,
Map inputsFromParent) {
- return inputsFromParent;
+ return Map.of("output", inputsFromParent.get("input"));
+ }
+
+ @Override
+ public ExecutorDescriptor getDescriptor(Map executorDescriptorParameters) {
+ return ExecutorDescriptor.builder().identifier("NoOp").name(name).inputName("input").outputName("output")
+ .description("No operation executor").build();
}
@Override
@@ -23,24 +36,19 @@ public class NoOpExecutor implements Executor {
return name;
}
- @Override
- public List getMandatoryParameters() {
- return List.of();
- }
-
- @Override
public Map> getRequiredPreExecutionParameters(Map userParameters) {
return Map.of();
}
- @Override
- public List declaredOutputNames() {
- return List.of();
- }
-
- @Override
- public List declaredInputNames() {
- return List.of();
+ public List getDynamicDescriptorParameters() {
+ ParameterDefinition parameterDefinition = ParameterDefinition.builder()
+ .name("IONumber")
+ .description("The Input/Output Number")
+ .type(ParameterType.Number).required(true)
+ .validation(Validations.minValidator(1))
+ .validation(Validations.maxValidator(4))
+ .build();
+ return List.of(parameterDefinition);
}
}
\ No newline at end of file
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/AIModel.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/AIModel.java
index c1f7f9b..d4a5c92 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/AIModel.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/AIModel.java
@@ -8,7 +8,7 @@ public interface AIModel {
String getDescription();
- String executePrompt(Map parameters, String prompt);
+ String executePrompt(Map parameters, String prompt) throws Throwable;
public Map> getRequiredPreExecutionParameters();
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GeminiModel.java b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GeminiModel.java
index 3218871..7501394 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GeminiModel.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executors/ai/GeminiModel.java
@@ -1,5 +1,6 @@
package it.cnr.isti.workflow.manager.executors.ai;
+import java.io.IOException;
import java.util.Map;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -8,6 +9,9 @@ import org.slf4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
import it.cnr.isti.workflow.manager.executors.ai.services.GeminiService;
@Service("google-gemini")
@@ -31,33 +35,27 @@ public class GeminiModel implements AIModel {
return "Gemini is an LLM model provided by Google";
}
-
@Override
- public String executePrompt(Map parameters, String prompt) {
+ public String executePrompt(Map parameters, String prompt) throws Throwable {
geminiService.setApiKey("***REMOVED-API-KEY***");
+ //log.info("------- REQUEST ------------");
+ //log.info(prompt);
+ //log.info("-------------------");
String result = geminiService.getResponse(prompt).block();
- log.debug("google gemini called with result {}",result);
+ ObjectMapper objectMapper = new ObjectMapper();
- //TODO: Remove, delay only for testing
- try {
- Thread.sleep(5000);
- } catch (InterruptedException e) {
- log.error("Error during sleep", e);
- }
+ // Deserializzazione della risposta JSON in un oggetto GeminiResponse
+ GeminiResponse response = objectMapper.readValue(result, GeminiResponse.class);
- result = parseResponse(result);
- log.debug("post parsing: {}",result);
- return result;
-
- }
+ // Estrazione del testo dal primo candidato
+ String candidateText = response.getCandidates().get(0).getContent().getParts().get(0).getText();
+
+ //log.info("------- RESPONSE ------------");
+ //log.info(candidateText);
+ //log.info("-------------------");
+
+ return candidateText;
- private String parseResponse(String response) {
- Pattern pattern = Pattern.compile("\"text\"\\s*:\\s*\"(.*?)\"");
- Matcher matcher = pattern.matcher(response);
- if (matcher.find()) {
- return matcher.group(1); // Estratto il testo
- }
- return response;
}
@@ -66,12 +64,56 @@ public class GeminiModel implements AIModel {
return Map.of("API-KEY", String.class);
}
- /*
- @Override
- public Map execute(Map userParameters, Map preExecutionParameters,
- Map inputsFromParent) {
- String returnString = (String)executePrompt(preExecutionParameters, (String) inputsFromParent.get("prompt"));
- return Map.of("response", returnString);
- }*/
-
+}
+
+@JsonIgnoreProperties(ignoreUnknown = true)
+class GeminiResponse {
+ private java.util.List candidates;
+
+ public java.util.List getCandidates() {
+ return candidates;
+ }
+
+ public void setCandidates(java.util.List candidates) {
+ this.candidates = candidates;
+ }
+
+ @JsonIgnoreProperties(ignoreUnknown = true)
+ public static class Candidate {
+ private Content content;
+
+ public Content getContent() {
+ return content;
+ }
+
+ public void setContent(Content content) {
+ this.content = content;
+ }
+ }
+
+ @JsonIgnoreProperties(ignoreUnknown = true)
+ public static class Content {
+ private java.util.List parts;
+
+ public java.util.List getParts() {
+ return parts;
+ }
+
+ public void setParts(java.util.List parts) {
+ this.parts = parts;
+ }
+ }
+
+ @JsonIgnoreProperties(ignoreUnknown = true)
+ public static class Part {
+ private String text;
+
+ public String getText() {
+ return text;
+ }
+
+ public void setText(String text) {
+ this.text = text;
+ }
+ }
}
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 81742a0..1d9a3df 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
@@ -3,14 +3,18 @@ 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.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
import it.cnr.isti.workflow.manager.executors.Executor;
+import it.cnr.isti.workflow.manager.model.ExecutorDescriptor;
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")
@@ -18,6 +22,7 @@ import lombok.ToString;
public class GenericAIExecutor implements Executor {
//EVERY executor in ai model must return a Map with the key "response" and the value of the response
+ private static final Logger log = LoggerFactory.getLogger(GenericAIExecutor.class);
private static final String name = "Generic AI Executor";
@@ -30,13 +35,18 @@ public class GenericAIExecutor implements Executor {
@Override
public Map execute(Map userParameters, Map preExecutionParameters,
- Map inputsFromParent) {
+ Map inputsFromParent) throws ExecutionException {
Objects.requireNonNull(userParameters.get("LLM"), "'LLM' parameter is required for 'GENERIC-AI' executor");
Objects.requireNonNull(aiModels.get(userParameters.get("LLM")), "AI model not found: " + userParameters.get("LLM"));
AIModel aiModel = aiModels.get(userParameters.get("LLM"));
String prompt = (String) inputsFromParent.get("prompt");
- String result = aiModel.executePrompt(userParameters, prompt);
- return Map.of("response", result);
+ try{
+ String result = aiModel.executePrompt(userParameters, prompt);
+ return Map.of("response", result);
+ }catch (Throwable e) {
+ log.error("Error executing AI model: {}", e.getMessage(), e);
+ throw new ExecutionException("Error executing AI model: " + e.getMessage(), e);
+ }
}
@@ -47,37 +57,29 @@ public class GenericAIExecutor implements Executor {
}
@Override
+ public ExecutorDescriptor getDescriptor(Map executorDescriptorParameters) {
+ return ExecutorDescriptor.builder().identifier("GENERIC-AI").name(name).inputName("prompt").outputName("response")
+ .description("Generic AI executor").mandatoryParameters(getMandatoryParameters()).build();
+ }
+
public Map> getRequiredPreExecutionParameters(Map userParameters) {
Objects.requireNonNull(userParameters.get("LLM"), "'LLM' parameter is required for 'GENERIC-AI' executor");
Objects.requireNonNull(aiModels.get(userParameters.get("LLM")), "AI model not found: " + userParameters.get("LLM"));
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())
+ .required(true)
.specificAttributes(rPars).build();
return List.of(rgParams);
}
-
-
-
+
}
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 987e00f..2da3633 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
@@ -4,7 +4,9 @@ import org.springframework.web.reactive.function.client.WebClient;
import org.springframework.http.MediaType;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
+import reactor.util.retry.Retry;
+import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -34,16 +36,27 @@ public class GeminiService {
)));
WebClient webClient = webClientBuilder.baseUrl(GEMINI_URL).build();
- Mono result = webClient.post()
+ Mono result = webClient.post()
.uri(uriBuilder -> uriBuilder.path("v1beta/models/gemini-2.0-flash:generateContent")
.queryParam("key", apiKey)
.build())
.contentType(MediaType.APPLICATION_JSON)
.bodyValue(requestBody)
.retrieve()
- .bodyToMono(String.class);
-
-
- return result;
+ .onStatus(
+ status -> status.is5xxServerError(),
+ clientResponse -> clientResponse.bodyToMono(String.class)
+ .defaultIfEmpty("Error: server without body")
+ .flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body))))
+ .bodyToMono(String.class)
+ .timeout(Duration.ofMinutes(2))
+ .retryWhen(
+ Retry.backoff(10, Duration.ofSeconds(30))
+ .filter(throwable -> throwable instanceof RuntimeException
+ || throwable instanceof java.util.concurrent.TimeoutException)
+ .onRetryExhaustedThrow((retryBackoffSpec, retrySignal) -> new RuntimeException(
+ "Error contacting Gemini API", retrySignal.failure())));
+
+ return result;
}
}
\ No newline at end of file
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 067dbf5..f793c47 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
@@ -9,9 +9,9 @@ import java.util.Set;
import org.slf4j.Logger;
+import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
-import it.cnr.isti.workflow.manager.model.types.Key;
import it.cnr.isti.workflow.manager.model.types.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.Translator;
import lombok.AllArgsConstructor;
@@ -37,13 +37,13 @@ public class ExecutionStep implements InputReadyListener {
boolean startAutomatically = true;
Executor executor;
- Map parentOutputsToInputMapping;
+ Map parentOutputsToInputMapping = new HashMap<>();
Map inputs = new HashMap<>();
List nextSteps = new ArrayList<>();
- List previouSteps = new ArrayList<>();
+ //List previouSteps = new ArrayList<>();
@NonNull
NodeDefinition nodeDefinition;
@@ -63,8 +63,6 @@ public class ExecutionStep implements InputReadyListener {
@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));
if (parentOutputsToInputMapping.containsKey(inputKey)) {
inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue);
if (startAutomatically) startExecution();
@@ -72,6 +70,7 @@ public class ExecutionStep implements InputReadyListener {
};
public boolean areAllInputsReady() {
+ log.debug("checking inputs {} - {} ({},{})", inputs.keySet(), parentOutputsToInputMapping.keySet(), inputs.size(), parentOutputsToInputMapping.size());
return inputs.size() == parentOutputsToInputMapping.size();
}
@@ -81,26 +80,22 @@ public class ExecutionStep implements InputReadyListener {
this.context.getStepsUnderExecution().add(this.id);
try {
Map realInputs = prepareInputs();
- log.debug("input translation for step {} is {}", id, realInputs);
+ //log.debug("input translation for step {} is {}", id, realInputs);
Map returned = executor.execute(this.runtimeParameters, null, realInputs);
- log.debug("returned from executor {} for step {} is {}", executor.getClass().getName(), id,
- returned);
this.context.getNodeResult().put(this.id, returned);
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) {
+ } 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.setStatus(Status.ERROR);
return;
} finally {
this.context.getStepsUnderExecution().remove(this.id);
- }
-
-
+ }
}).start();
else
log.debug("not all inputs are ready for step {}", id);
@@ -111,23 +106,21 @@ public class ExecutionStep implements InputReadyListener {
Map realInputs = new HashMap<>();
Set keys = new HashSet<>(inputs.keySet());
if (nodeDefinition.getInputTranslators() != null) {
- log.debug("translators {}", nodeDefinition.getInputTranslators().size());
+ //log.debug("translators {}", nodeDefinition.getInputTranslators().size());
for (Map.Entry entry : nodeDefinition.getInputTranslators().entrySet()) {
String translation = entry.getValue().getTranslation();
-
- for (Key usedKey : entry.getValue().getUsedKeys()) {
+ List usedKeys = entry.getValue().getUsedKeys();
+ log.debug("used keys", usedKeys.toString());
+ for (String usedKey : usedKeys) {
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);
+ if (this.runtimeParameters.containsKey(usedKey))
+ valueToReplace = (String) this.runtimeParameters.get(usedKey);
+ else if (inputs.containsKey(usedKey)){
+ valueToReplace = (String) inputs.get(usedKey);
+ keys.remove(usedKey);
+ }
+ translation = translation.replace(String.format("${{%s}}", usedKey), valueToReplace);
+ //log.debug("replacing used key {} value {}", usedKey, valueToReplace);
}
realInputs.put(entry.getKey(), translation);
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java
index 95c437b..305e123 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/model/ExecutorDescriptor.java
@@ -4,27 +4,33 @@ import java.util.List;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
import lombok.Data;
+import lombok.NonNull;
+import lombok.Singular;
@Data
+@AllArgsConstructor
+@Builder
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();
- }
-
+ @NonNull
String identifier;
+ @NonNull
String name;
+
String description;
+ @Singular
List inputNames;
+ @Singular
List outputNames;
+ @Singular
List mandatoryParameters;
+ List dynamicDescriptorParameter;
+
//List optionalParameters;
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/Validations.java b/src/main/java/it/cnr/isti/workflow/manager/model/Validations.java
new file mode 100644
index 0000000..cca68ba
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/model/Validations.java
@@ -0,0 +1,36 @@
+package it.cnr.isti.workflow.manager.model;
+
+
+import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
+import it.cnr.isti.workflow.manager.model.types.Validation;
+
+public class Validations {
+
+
+ public static Validation maxValidator(int value){
+ return Validation.builder().name("max").value(value).message("max value is " + value).with((o) -> ((Integer) o)<=value ).build();
+ }
+ public static Validation minValidator(int value){
+ return Validation.builder().name("min").value(value).message("min value is " + value).with((o) -> ((Integer) o)>=value).build();
+ }
+ public static Validation maxLengthValidator(int value){
+ return Validation.builder().name("maxLength").value(value).message("max length is " + value).with((o) -> o.toString().length() o.toString().length()>value).build();
+ }
+
+ public static Validation regexValidator(String regex){
+ return Validation.builder().name("pattern").value(regex).message("regex is " + regex).with((o) -> o.toString().matches(regex)).build();
+ }
+
+
+ public static boolean isValid(ParameterDefinition parameter, Object value) {
+ for (Validation validation : parameter.getValidations()) {
+ if (!validation.getWith().test(value)) {
+ return false;
+ }
+ }
+ return true;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/IOType.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/IOType.java
index 1f03f2f..9ed2d2b 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/model/types/IOType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/IOType.java
@@ -1,6 +1,7 @@
package it.cnr.isti.workflow.manager.model.types;
import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonValue;
public enum IOType {
Text,
@@ -10,4 +11,9 @@ public enum IOType {
public static IOType fromString(String key) {
return IOType.valueOf(key);
}
+
+ @JsonValue
+ public String toString() {
+ return name();
+ }
}
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 58c7fc0..39b44b6 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
@@ -36,6 +36,9 @@ public class NodeDefinition {
@NonNull
private String executor;
+ @Builder.Default
+ private String color = "black";
+
@Singular
@Convert(converter = MapsConverter.class)
@Column(columnDefinition = "TEXT")
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 a92775f..dc89061 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
@@ -38,6 +38,9 @@ public class ParameterDefinition {
@NonNull
private ParameterType type;
+ @Builder.Default
+ private boolean required = false;
+
@Singular
private List validations;
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 5b40713..a4d1605 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
@@ -3,5 +3,7 @@ package it.cnr.isti.workflow.manager.model.types;
public enum ParameterType {
Select,
- Input
+ Text,
+ Boolean,
+ Number
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/Translator.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/Translator.java
index 6245bf3..ec9b1cb 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/model/types/Translator.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/Translator.java
@@ -1,13 +1,18 @@
package it.cnr.isti.workflow.manager.model.types;
+import java.util.ArrayList;
import java.util.List;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+import org.slf4j.Logger;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.NonNull;
-import lombok.Singular;
@Data
@@ -16,17 +21,32 @@ import lombok.Singular;
@Builder
public class Translator {
- @Singular
- @NonNull
- List usedKeys;
+ private static final Logger log = org.slf4j.LoggerFactory.getLogger(Translator.class);
+
+ private static final Pattern pattern = Pattern.compile("\\$\\{\\{(.*?)\\}\\}");
@NonNull
String translation;
public String toString() {
- return "Translator(usedKeys=" + this.getUsedKeys() + ", translation=" + this.getTranslation() + ")";
+ return "Translator(translation=" + this.getTranslation() + ")";
}
+ @JsonIgnore
+ public List getUsedKeys() {
+ Matcher matcher = pattern.matcher(translation);
+
+ List variables = new ArrayList<>();
+
+ // Cicla tutti i match trovati
+ while (matcher.find()) {
+ String found =matcher.group(1);
+ log.info("found: {} ",found);
+ variables.add(found); // prende solo il nome della variabile (senza ${{}})
+ }
+
+ return variables;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/model/types/Validation.java b/src/main/java/it/cnr/isti/workflow/manager/model/types/Validation.java
index 372f122..98dbf48 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/model/types/Validation.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/model/types/Validation.java
@@ -1,8 +1,11 @@
package it.cnr.isti.workflow.manager.model.types;
+import java.util.function.Predicate;
+
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
+import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import lombok.NonNull;
@@ -10,6 +13,7 @@ import lombok.NonNull;
@AllArgsConstructor
@Data
@NoArgsConstructor
+@EqualsAndHashCode
public class Validation {
@NonNull
@@ -18,5 +22,12 @@ public class Validation {
@NonNull
private String validator;
+ private Object value;
+
+ @EqualsAndHashCode.Exclude
private String message;
+
+ @NonNull
+ @EqualsAndHashCode.Exclude
+ private Predicate