HumanInteraction node added

This commit is contained in:
Lucio Lelii 2025-07-30 16:59:17 +02:00
parent 4200a2bc3d
commit 985d2fb67f
39 changed files with 797 additions and 329 deletions

View File

@ -7,13 +7,13 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Import;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.OutputNodeDefinition;
import it.cnr.isti.workflow.manager.executors.ai.models.AIModel;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import it.cnr.isti.workflow.manager.services.ImportComponent;
@ -36,14 +36,20 @@ public class WorkflowManagerApplication {
return args -> {
NodeDefinition textInput = InputNodeDefinition.builder().name("Text Input")
.createdBy("admin").category("Inputs").outputType(IOType.TEXT)
.outputType(IOType.TEXT)
.build();
NodeDefinition textOutput = OutputNodeDefinition.builder().name("Text Output")
.createdBy("admin").category("Outputs").inputType(IOType.TEXT)
.inputType(IOType.TEXT)
.build();
HumanInteractionNodeDefinition humanInteractionNodeDefinition = HumanInteractionNodeDefinition.builder()
.name("Human Interaction")
.simulable(false)
.build();
repository.save(textInput);
repository.save(textOutput);
repository.save(humanInteractionNodeDefinition);
importComponent.start();
};

View File

@ -90,9 +90,8 @@ public class FlowsController {
if (flow.getId() != null)
throw new WebException(HttpStatus.BAD_REQUEST, "on creation flow id must be null");
flow.setCreatedBy(userDetails.getUsername());
Flow savedFlow = flowRepository.save((FlowEntity)flow);
logger.info("flow created {}", flow.getId());
FlowEntity savedFlow = flowRepository.save(FlowEntity.from(flow));
logger.info("flow created {}", savedFlow.getId());
return savedFlow;
}
@ -113,7 +112,7 @@ public class FlowsController {
//validateFlow(flow);
if (!id.equals(flow.getId()))
throw new WebException(HttpStatus.BAD_REQUEST, "Flow id does not match");
Flow savedFlow = flowRepository.findById(id).orElseThrow(() -> new WebException(HttpStatus.NOT_FOUND,"Flow not found"));
FlowEntity savedFlow = flowRepository.findById(id).orElseThrow(() -> new WebException(HttpStatus.NOT_FOUND,"Flow not found"));
if (!savedFlow.getCreatedBy().equals(userDetails.getUsername()))
throw new WebException(HttpStatus.FORBIDDEN, "Flow not created by user");

View File

@ -1,4 +1,4 @@
package it.cnr.isti.workflow.manager.controllers.nodes;
package it.cnr.isti.workflow.manager.controllers;
import java.util.List;
@ -11,11 +11,13 @@ import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutorDescriptor;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import it.cnr.isti.workflow.manager.repositories.SystemNodeDefinitionRepository;
import it.cnr.isti.workflow.manager.repositories.UserNodeDefinitionRepository;
@RestController
@RequestMapping("/types")
@ -26,26 +28,44 @@ public class NodeDefinitionController {
@Autowired
private NodeDefinitionRepository repository;
@Autowired
private UserNodeDefinitionRepository userNodeDefinitionRepository;
@Autowired
private SystemNodeDefinitionRepository systemNodeDefinitionRepository;
@Autowired
private List<Executor> executors;
//@SecurityRequirement(name = "bearerAuth")
@GetMapping("/nodes")
public List<NodeDefinition> getNodes() {
return (List<NodeDefinition>) repository.findAll();
@GetMapping("/nodes/system")
public List<SystemNodeDefinition> getSystemNodes() {
return (List<SystemNodeDefinition>) systemNodeDefinitionRepository.findAll();
}
@GetMapping("/nodes/types")
public List<String> getStoredNodeTypes() {
return repository.getUsedTypes();
@GetMapping("/nodes/user")
public List<UserNodeDefinition> getUserNodes() {
return (List<UserNodeDefinition>) userNodeDefinitionRepository.findAll();
}
@GetMapping("/nodes/system/categories")
public List<String> getSystemNodeCategories() {
return systemNodeDefinitionRepository.getCategories();
}
@GetMapping("/nodes/user/types")
public List<String> getUserNodeTypes() {
return userNodeDefinitionRepository.getUsedTypes();
}
@GetMapping("/nodes/categories")
public List<String> getStoredNodeCategories() {
return repository.getCategories();
@GetMapping("/nodes/user/categories")
public List<String> getUserNodeCategories() {
return userNodeDefinitionRepository.getUsedCategories();
}
@SecurityRequirement(name = "bearerAuth")

View File

@ -2,39 +2,61 @@ package it.cnr.isti.workflow.manager.dto;
import java.util.List;
import com.fasterxml.jackson.annotation.JsonAlias;
import lombok.AllArgsConstructor;
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import lombok.AccessLevel;
import lombok.Builder;
import lombok.Builder.Default;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.NonNull;
import lombok.ToString;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
@NoArgsConstructor(access = AccessLevel.PROTECTED)
@ToString
public class Flow {
@Builder
protected Flow(String id, String createdBy, boolean isPublic, String name, String description, List<Node> nodes,
List<Connection> connections) {
this.id = id;
this.createdBy = createdBy;
this.isPublic = isPublic;
this.name = name;
this.description = description;
this.nodes = nodes;
this.connections = connections;
if (nodes != null) {
this.loadNodeDefinitions();
}
}
String id;
@NonNull
String createdBy;
@Default
@JsonAlias("public")
boolean isPublic =false;
boolean isPublic = false;
@NonNull
String name;
String description;
@Builder.Default
List<Node> nodes = List.of();
protected List<Node> nodes = List.of();
public void setNodes(List<Node> nodes) {
this.nodes = nodes;
if (nodes != null) {
this.loadNodeDefinitions();
}
}
@Builder.Default
List<Connection> connections = List.of();
protected void loadNodeDefinitions() {
this.getNodes().forEach(Node::resolveNodeDefinition);
}
}

View File

@ -3,10 +3,9 @@ package it.cnr.isti.workflow.manager.dto;
import java.util.List;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.services.RepositoryHolder;
import lombok.AllArgsConstructor;
import lombok.Builder;
@ -46,7 +45,7 @@ public class Node {
String type;
public void resolveNodeDefinition() {
if (type != null) {
if (type != null && nodeDefinition == null) {
this.nodeDefinition = RepositoryHolder.getNodeDefinition()
.findById(type)
.orElse(null);

View File

@ -14,6 +14,8 @@ import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
import jakarta.persistence.Id;
import jakarta.persistence.PostLoad;
import jakarta.persistence.PostPersist;
import jakarta.persistence.PostUpdate;
import lombok.EqualsAndHashCode;
@Entity
@ -106,7 +108,7 @@ public class FlowEntity extends Flow {
@Override
public void setNodes(List<Node> nodes) {
super.setNodes(nodes);
this.nodes = nodes;
}
@Convert(converter = ConnectionListConverter.class)
@ -121,8 +123,10 @@ public class FlowEntity extends Flow {
super.setConnections(connections);
}
@PostPersist
@PostUpdate
@PostLoad
public void postLoad() {
this.getNodes().forEach(Node::resolveNodeDefinition);
super.loadNodeDefinitions();
}
}

View File

@ -0,0 +1,40 @@
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
import it.cnr.isti.workflow.manager.model.types.ParameterType;
import jakarta.persistence.Entity;
import lombok.Builder;
import lombok.NoArgsConstructor;
@Entity
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@JsonTypeName("system")
public class HumanInteractionNodeDefinition extends SystemNodeDefinition {
public static final String PARAMETER_NAME = "action_description";
public static final String INPUT_NAME = "input";
public static final String OUTPUT_NAME = "output";
@Builder
public HumanInteractionNodeDefinition(String name, boolean simulable) {
super(name, simulable);
this.setCategory("humanInteraction");
this.outputs = Map.of(OUTPUT_NAME, IOType.TEXT);
this.inputs = Map.of(INPUT_NAME, IOType.TEXT);
this.setColor("black");
this.runtimeParameters = Map.of(PARAMETER_NAME, ParameterDefinition.builder()
.name(PARAMETER_NAME)
.label("Human Readable Action Description")
.description("Description of the user action to perform")
.required(true)
.type(ParameterType.Text)
.build());
}
}

View File

@ -1,37 +1,40 @@
package it.cnr.isti.workflow.manager.model.types.nodes.definitions;
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
import it.cnr.isti.workflow.manager.model.types.ParameterType;
import jakarta.persistence.DiscriminatorValue;
import jakarta.persistence.Entity;
import lombok.Builder;
import lombok.NoArgsConstructor;
@Entity
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@DiscriminatorValue("2")
public class InputNodeDefinition extends NodeDefinition {
@JsonTypeName("system")
public class InputNodeDefinition extends SystemNodeDefinition {
public static final String INPUT_PARAMETER_KEY ="name";
public static final String OUTPUT_NAME = "value";
public static final String DEFAULT_VALUE = "output";
@Builder
public InputNodeDefinition(String name, String createdBy, String category, IOType outputType) {
public InputNodeDefinition(String name, IOType outputType, boolean simulable) {
super();
this.simulable = simulable;
this.setName(name);
this.setCreatedBy(createdBy);
this.setCategory(category);
this.setCategory("input");
this.outputs = Map.of(OUTPUT_NAME, outputType);
this.setColor("blue");
this.setColor("green");
this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder()
.name(INPUT_PARAMETER_KEY)
.label("Input Name")
.description("Input name")
.defaultValue("output")
.defaultValue(DEFAULT_VALUE)
.type(ParameterType.String)
.build());
}

View File

@ -1,8 +1,10 @@
package it.cnr.isti.workflow.manager.model.types.nodes.definitions;
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import java.util.HashMap;
import java.util.Map;
import org.springframework.boot.autoconfigure.security.SecurityProperties.User;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@ -35,13 +37,7 @@ import lombok.experimental.SuperBuilder;
@Inheritance(strategy = InheritanceType.SINGLE_TABLE)
@DiscriminatorColumn(name="node_type", length = 1,
discriminatorType = DiscriminatorType.INTEGER)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type", defaultImpl = UserNodeDefinition.class)
@JsonSubTypes({
@JsonSubTypes.Type(value = InputNodeDefinition.class, name = "input"),
@JsonSubTypes.Type(value = UserNodeDefinition.class, name = "user"),
@JsonSubTypes.Type(value = OutputNodeDefinition.class, name = "output"),
})
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type", defaultImpl = UserNodeDefinition.class)
public abstract class NodeDefinition {
@Id

View File

@ -1,34 +1,33 @@
package it.cnr.isti.workflow.manager.model.types.nodes.definitions;
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import java.util.Map;
import it.cnr.isti.workflow.manager.model.Validations;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
import it.cnr.isti.workflow.manager.model.types.ParameterType;
import it.cnr.isti.workflow.manager.model.types.Validation;
import jakarta.persistence.DiscriminatorValue;
import jakarta.persistence.Entity;
import lombok.Builder;
import lombok.NoArgsConstructor;
@Entity
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@DiscriminatorValue("3")
public class OutputNodeDefinition extends NodeDefinition {
@JsonTypeName("system")
public class OutputNodeDefinition extends SystemNodeDefinition {
public static final String INPUT_PARAMETER_KEY ="name";
public static final String INPUT_NAME = "value";
@Builder
public OutputNodeDefinition(String name, String createdBy, String category, IOType inputType) {
public OutputNodeDefinition(String name, IOType inputType, boolean simulable) {
super();
this.simulable = simulable;
this.setName(name);
this.setCreatedBy(createdBy);
this.setCategory(category);
this.setCategory("output");
this.inputs = Map.of(INPUT_NAME, inputType);
this.setColor("blue");
this.setColor("red");
this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder()
.name(INPUT_PARAMETER_KEY)
.label("Output Name")

View File

@ -0,0 +1,25 @@
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.persistence.DiscriminatorValue;
import jakarta.persistence.Entity;
import lombok.NoArgsConstructor;
@Entity
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@DiscriminatorValue("1")
public abstract class SystemNodeDefinition extends NodeDefinition {
@JsonProperty
boolean simulable = false;
public SystemNodeDefinition(String name, boolean simulable) {
super();
this.simulable = simulable;
this.setName(name);
this.setCreatedBy("system");
}
}

View File

@ -1,6 +1,10 @@
package it.cnr.isti.workflow.manager.model.types.nodes.definitions;
package it.cnr.isti.workflow.manager.entities.nodes.definitions;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.model.types.Translator;
import it.cnr.isti.workflow.manager.repositories.converters.InputTranslatorConverter;
import it.cnr.isti.workflow.manager.repositories.converters.MapsConverter;
@ -23,7 +27,8 @@ import lombok.experimental.SuperBuilder;
@SuperBuilder
@Data
@EqualsAndHashCode(callSuper = true)
@DiscriminatorValue("1")
@DiscriminatorValue("2")
@JsonTypeName("user")
public class UserNodeDefinition extends NodeDefinition {
@NonNull

View File

@ -5,9 +5,6 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.FieldKey;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;

View File

@ -1,8 +1,6 @@
package it.cnr.isti.workflow.manager.executors;
import org.springframework.stereotype.Service;
import java.lang.reflect.Parameter;
import java.util.List;
import java.util.Map;

View File

@ -1,7 +1,6 @@
package it.cnr.isti.workflow.manager.executors.ai.models;
import java.util.Map;
import org.slf4j.Logger;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
@ -11,7 +10,7 @@ import it.cnr.isti.workflow.manager.executors.ai.services.GeminiService;
@Service("google-gemini")
public class GeminiModel implements AIModel {
private static final Logger log = org.slf4j.LoggerFactory.getLogger(GeminiModel.class);
//private static final Logger log = org.slf4j.LoggerFactory.getLogger(GeminiModel.class);
private GeminiService geminiService;

View File

@ -13,20 +13,23 @@ import lombok.NoArgsConstructor;
public class ExecutionContext {
public enum Status {
CREATED(true, false),
INITIALIZING(true, false),
READY(true, false),
RUNNING(false, false),
SUCCESS(false, true),
ERROR(false, true);
CREATED(true, false, false),
INITIALIZING(true, false, false),
READY(true, false, false),
RUNNING(false, false, false),
WAITING(false, false, true),
SUCCESS(false, true, false),
ERROR(false, true, false);
private boolean initState;
private boolean finalState;
private boolean waitingState;
Status(boolean initState, boolean finalState) {
Status(boolean initState, boolean finalState, boolean waitingState) {
this.initState = initState;
this.finalState = finalState;
this.waitingState = waitingState;
}
public boolean isInitState() {
@ -37,8 +40,12 @@ public class ExecutionContext {
return finalState;
}
public boolean isWaitingState() {
return waitingState ;
}
public boolean isRunningState() {
return !finalState && !initState;
return !finalState && !initState && !waitingState;
}
}
@ -51,6 +58,7 @@ public class ExecutionContext {
Long endTime = null;
List<String> stepsUnderExecution = new ArrayList<>();
List<String> waitingSteps = new ArrayList<>();
Map<String, String> errors = new HashMap<>();
List<String> warnings = new ArrayList<>();

View File

@ -2,7 +2,6 @@ package it.cnr.isti.workflow.manager.model;
import java.util.List;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
import lombok.AllArgsConstructor;
import lombok.Builder;

View File

@ -1,10 +1,8 @@
package it.cnr.isti.workflow.manager.model;
import java.util.Map;
public interface InputReadyListener {
void inputReady(String key, Object value);
Map<String, String> getParentOutputsToInputMapping();
void addParentOutputsToInputMapping(String parentOutputName, String inputName);
}

View File

@ -1,165 +1,12 @@
package it.cnr.isti.workflow.manager.model.steps;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import it.cnr.isti.workflow.manager.dto.Point;
import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.Translator;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import lombok.Builder;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.NonNull;
import lombok.ToString;
@Data
@ToString
@EqualsAndHashCode(callSuper = true)
public class ExecutionStep extends Step implements InputReadyListener, ExecutableStep {
public abstract class ExecutionStep extends Step {
private static final Logger log = org.slf4j.LoggerFactory.getLogger(ExecutionStep.class);
@NonNull
String name;
Map<String, Object> runtimeParameters;
Map<String, Object> fixedParameters;
@JsonIgnore
@NonNull
ExecutionContext context;
Map<String, IOType> inputTypes;
Map<String, IOType> outputTypes;
@JsonIgnore
Executor executor;
@JsonIgnore
Map<String, String> parentOutputsToInputMapping = new HashMap<>();
@JsonIgnore
Map<String, Object> inputs = new HashMap<>();
@JsonIgnore
@NonNull
UserNodeDefinition nodeDefinition;
@Builder
public ExecutionStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map<String, Object> runtimeParameters,
Map<String, Object> fixedParameters,
UserNodeDefinition nodeDefinition,
Executor executor, Point position, @NonNull Map<String, IOType> inputTypes,
@NonNull Map<String, IOType> outputTypes) {
public ExecutionStep(String id, Point position) {
super(id, position);
this.name = name;
this.runtimeParameters = runtimeParameters;
this.fixedParameters = fixedParameters;
this.executor = executor;
this.nodeDefinition = nodeDefinition;
this.context = context;
this.inputTypes = inputTypes;
this.outputTypes = outputTypes;
}
@Override
public void inputReady(String inputKey, Object inputValue) {
if (parentOutputsToInputMapping.containsKey(inputKey)) {
inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue);
if (startAutomatically) start();
}
};
public boolean areAllInputsReady() {
log.debug("checking inputs {} - {} ({},{})", inputs.keySet(), parentOutputsToInputMapping.keySet(), inputs.size(), parentOutputsToInputMapping.size());
return inputs.size() == parentOutputsToInputMapping.size();
}
public void start() {
if (areAllInputsReady())
new Thread(() -> {
this.context.getStepsUnderExecution().add(this.id);
log.info("starting executing of step {}", this.nodeDefinition.getName());
try {
Map<String, Object> realInputs = prepareInputs();
//log.debug("input translation for step {} is {}", id, realInputs);
Map<String, Object> returned = executor.execute(this.runtimeParameters, null, realInputs);
//mapping executor output to current node output
Map<String, Object> mappedReturn = returned.entrySet().stream().collect(Collectors.toMap(
e -> { return !nodeDefinition.getExecutorToOutputTranslationMappings().isEmpty()
? nodeDefinition.getExecutorToOutputTranslationMappings().get(e.getKey())
: e.getKey();
},
e -> e.getValue()));
context.addNodeResult(id, mappedReturn);
//sendig output to next steps
nextSteps.forEach(n -> mappedReturn.forEach((k, v) -> n.inputReady(k,v)));
} catch (ExecutionException e) {
log.error("error during execution of step {}: {}", id, e.getMessage());
this.context.addError(this.id, e.getMessage());
this.context.setStatus(Status.ERROR);
return;
} finally {
this.context.getStepsUnderExecution().remove(this.id);
log.info("finished executing step {}", this.nodeDefinition.getName());
}
}).start();
else
log.debug("not all inputs are ready for step {}", id);
}
private Map<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();
List<String> usedKeys = entry.getValue().getUsedKeys();
for (String usedKey : usedKeys) {
String valueToReplace = "";
if (this.runtimeParameters.containsKey(usedKey))
valueToReplace = (String) this.runtimeParameters.get(usedKey);
else if (inputs.containsKey(usedKey)){
valueToReplace = (String) inputs.get(usedKey);
keys.remove(usedKey);
}
translation = translation.replace(String.format("${{%s}}", usedKey), valueToReplace);
//log.debug("replacing used key {} value {}", usedKey, valueToReplace);
}
realInputs.put(entry.getKey(), translation);
}
}
// key is the executor input name, value is the input name
if (this.nodeDefinition.getInputToExecutorTranslationMappings() != null) {
for (Map.Entry<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;
}
}

View File

@ -1,9 +1,10 @@
package it.cnr.isti.workflow.manager.model.steps;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.dto.Point;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import lombok.Builder;
import lombok.Data;
import lombok.EqualsAndHashCode;
@ -11,6 +12,7 @@ import lombok.NonNull;
@Data
@EqualsAndHashCode(callSuper = true)
@JsonTypeName("input")
public class InputStep extends Step implements ExecutableStep {
IOType inputType;

View File

@ -0,0 +1,78 @@
package it.cnr.isti.workflow.manager.model.steps;
import java.util.HashMap;
import java.util.Map;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.dto.Point;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.types.IOType;
import lombok.Builder;
import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.NonNull;
import lombok.ToString;
@ToString
@Getter
@EqualsAndHashCode(callSuper = true)
@JsonTypeName("interactive")
public class InteractiveStep extends ExecutionStep implements InputReadyListener {
private static Logger log = LoggerFactory.getLogger(InteractiveStep.class);
IOType inputType;
IOType outputType;
String interactionDescription;
@JsonIgnore
ExecutionContext context;
Map<String, String> parentOutputsToInputMapping = new HashMap<>();
Object inputValue;
@Builder
public InteractiveStep(@NonNull String id, @NonNull ExecutionContext context,
@NonNull IOType inputType,
@NonNull IOType outputType, @NonNull String interactionDescription,
Point position) {
super(id, position);
this.interactionDescription = interactionDescription;
this.inputType = inputType;
this.outputType = outputType;
this.context = context;
}
@Override
public void inputReady(String key, Object value) {
log.debug("Input ready for Interactive step: {}, key: {}, value: {}", this.getId(), key, value);
this.inputValue = value;
this.context.setStatus(Status.WAITING);
this.context.getWaitingSteps().add(this.id);
}
@Override
public void addParentOutputsToInputMapping(String parentOutputName, String inputName) {
this.parentOutputsToInputMapping.put(parentOutputName, inputName);
}
public void setUserInput(Object input) {
this.inputValue = input;
this.context.setStatus(Status.RUNNING);
this.context.getWaitingSteps().remove(this.id);
this.getNextSteps().forEach(step -> {
step.inputReady(HumanInteractionNodeDefinition.OUTPUT_NAME, inputValue);
});
}
}

View File

@ -7,6 +7,7 @@ import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonTypeName;
import io.micrometer.common.lang.NonNull;
import it.cnr.isti.workflow.manager.dto.Point;
@ -18,6 +19,7 @@ import lombok.EqualsAndHashCode;
@Data
@EqualsAndHashCode(callSuper = true)
@JsonTypeName("output")
public class OutputStep extends Step implements InputReadyListener, ResultProvider {
private IOType outputType;
@ -54,5 +56,10 @@ public class OutputStep extends Step implements InputReadyListener, ResultProvid
this.resultListeners.forEach(listener -> listener.resultReady(this.id, this.getOutputName(), value));
}
@Override
public void addParentOutputsToInputMapping(String parentOutputName, String inputName) {
this.parentOutputsToInputMapping.put(parentOutputName, inputName);
}
}

View File

@ -7,7 +7,6 @@ import java.util.Map;
import it.cnr.isti.workflow.manager.dto.Point;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import io.micrometer.common.lang.NonNull;
@ -18,11 +17,6 @@ import lombok.Data;
@Data
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "type")
@JsonSubTypes({
@JsonSubTypes.Type(value = InputStep.class, name = "input"),
@JsonSubTypes.Type(value = OutputStep.class, name = "output"),
@JsonSubTypes.Type(value = ExecutionStep.class, name = "execution")
})
public abstract class Step {
@NonNull

View File

@ -0,0 +1,171 @@
package it.cnr.isti.workflow.manager.model.steps;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import it.cnr.isti.workflow.manager.dto.Point;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonTypeName;
import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.Translator;
import lombok.Builder;
import lombok.EqualsAndHashCode;
import lombok.NonNull;
import lombok.ToString;
@ToString
@EqualsAndHashCode(callSuper = true)
@JsonTypeName("execution")
public class UserDefinedStep extends ExecutionStep implements InputReadyListener, ExecutableStep {
private static final Logger log = org.slf4j.LoggerFactory.getLogger(UserDefinedStep.class);
@NonNull
String name;
Map<String, Object> runtimeParameters;
Map<String, Object> fixedParameters;
@JsonIgnore
@NonNull
ExecutionContext context;
Map<String, IOType> inputTypes;
Map<String, IOType> outputTypes;
@JsonIgnore
Executor executor;
Map<String, String> parentOutputsToInputMapping = new HashMap<>();
@JsonIgnore
Map<String, Object> inputs = new HashMap<>();
@JsonIgnore
@NonNull
UserNodeDefinition nodeDefinition;
@Builder
public UserDefinedStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map<String, Object> runtimeParameters,
Map<String, Object> fixedParameters,
UserNodeDefinition nodeDefinition,
Executor executor, Point position, @NonNull Map<String, IOType> inputTypes,
@NonNull Map<String, IOType> outputTypes) {
super(id, position);
this.name = name;
this.runtimeParameters = runtimeParameters;
this.fixedParameters = fixedParameters;
this.executor = executor;
this.nodeDefinition = nodeDefinition;
this.context = context;
this.inputTypes = inputTypes;
this.outputTypes = outputTypes;
}
@Override
public void inputReady(String inputKey, Object inputValue) {
if (parentOutputsToInputMapping.containsKey(inputKey)) {
inputs.put(parentOutputsToInputMapping.get(inputKey), inputValue);
if (startAutomatically) start();
}
};
public boolean areAllInputsReady() {
log.debug("checking inputs {} - {} ({},{})", inputs.keySet(), parentOutputsToInputMapping.keySet(), inputs.size(), parentOutputsToInputMapping.size());
return inputs.size() == parentOutputsToInputMapping.size();
}
public void start() {
if (areAllInputsReady())
new Thread(() -> {
this.context.getStepsUnderExecution().add(this.id);
log.info("starting executing of step {}", this.nodeDefinition.getName());
try {
Map<String, Object> realInputs = prepareInputs();
//log.debug("input translation for step {} is {}", id, realInputs);
Map<String, Object> returned = executor.execute(this.runtimeParameters, null, realInputs);
//mapping executor output to current node output
Map<String, Object> mappedReturn = returned.entrySet().stream().collect(Collectors.toMap(
e -> { return !nodeDefinition.getExecutorToOutputTranslationMappings().isEmpty()
? nodeDefinition.getExecutorToOutputTranslationMappings().get(e.getKey())
: e.getKey();
},
e -> e.getValue()));
context.addNodeResult(id, mappedReturn);
//sendig output to next steps
nextSteps.forEach(n -> mappedReturn.forEach((k, v) -> n.inputReady(k,v)));
} catch (ExecutionException e) {
log.error("error during execution of step {}: {}", id, e.getMessage());
this.context.addError(this.id, e.getMessage());
this.context.setStatus(Status.ERROR);
return;
} finally {
this.context.getStepsUnderExecution().remove(this.id);
log.info("finished executing step {}", this.nodeDefinition.getName());
}
}).start();
else
log.debug("not all inputs are ready for step {}", id);
}
private Map<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();
List<String> usedKeys = entry.getValue().getUsedKeys();
for (String usedKey : usedKeys) {
String valueToReplace = "";
if (this.runtimeParameters.containsKey(usedKey))
valueToReplace = (String) this.runtimeParameters.get(usedKey);
else if (inputs.containsKey(usedKey)){
valueToReplace = (String) inputs.get(usedKey);
keys.remove(usedKey);
}
translation = translation.replace(String.format("${{%s}}", usedKey), valueToReplace);
//log.debug("replacing used key {} value {}", usedKey, valueToReplace);
}
realInputs.put(entry.getKey(), translation);
}
}
// key is the executor input name, value is the input name
if (this.nodeDefinition.getInputToExecutorTranslationMappings() != null) {
for (Map.Entry<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;
}
@Override
public void addParentOutputsToInputMapping(String parentOutputName, String inputName) {
this.parentOutputsToInputMapping.put(parentOutputName, inputName);
}
}

View File

@ -3,7 +3,6 @@ package it.cnr.isti.workflow.manager.model.types;
import java.io.IOException;
import java.util.List;
import org.apache.commons.lang3.Validate;
import org.json.JSONObject;
import com.fasterxml.jackson.core.JacksonException;

View File

@ -1,10 +0,0 @@
package it.cnr.isti.workflow.manager.repositories;
import org.springframework.data.repository.CrudRepository;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
public interface InputDefinitionRepository extends CrudRepository<InputNodeDefinition, String> {
}

View File

@ -1,19 +1,12 @@
package it.cnr.isti.workflow.manager.repositories;
import java.util.List;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.CrudRepository;
import org.springframework.stereotype.Repository;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
@Repository
public interface NodeDefinitionRepository extends CrudRepository<NodeDefinition, String> {
@Query("SELECT n.name FROM NodeDefinition n")
List<String> getUsedTypes();
@Query("SELECT DISTINCT n.category FROM NodeDefinition n")
List<String> getCategories();
}

View File

@ -1,9 +0,0 @@
package it.cnr.isti.workflow.manager.repositories;
import org.springframework.data.repository.CrudRepository;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition;
public interface OutputDefinitionRepository extends CrudRepository<OutputNodeDefinition,String> {
}

View File

@ -0,0 +1,15 @@
package it.cnr.isti.workflow.manager.repositories;
import java.util.List;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.CrudRepository;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition;
public interface SystemNodeDefinitionRepository extends CrudRepository<SystemNodeDefinition,String> {
@Query("SELECT DISTINCT n.category FROM SystemNodeDefinition n")
List<String> getCategories();
}

View File

@ -1,9 +1,17 @@
package it.cnr.isti.workflow.manager.repositories;
import java.util.List;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.CrudRepository;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
public interface UserNodeDefinitionRepository extends CrudRepository<UserNodeDefinition, String> {
@Query("SELECT n.name FROM UserNodeDefinition n")
List<String> getUsedTypes();
@Query("SELECT DISTINCT n.category FROM UserNodeDefinition n")
List<String> getUsedCategories();
}

View File

@ -8,7 +8,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.steps.ExecutionStep;
import it.cnr.isti.workflow.manager.model.steps.InputStep;
import it.cnr.isti.workflow.manager.model.steps.InteractiveStep;
@Service
public class ExecutionService {
@ -39,45 +41,59 @@ public class ExecutionService {
if (!executions.containsKey(id))
throw new IllegalArgumentException("Execution with id " + id + " not found");
if (executions.get(id).getContext().getStatus() == Status.RUNNING)
throw new IllegalStateException("Execution with id "+id+" is still running");
throw new IllegalStateException("Execution with id " + id + " is still running");
executions.remove(id);
}
public ExecutionObject prepareInput(String executionId, String nodeId, String inputName, Object input) {
ExecutionObject eo = getExecution(executionId);
if (!eo.getContext().getStatus().isInitState()) {
if (eo.getContext().getStatus().isInitState()) {
InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found"));
if (!step.getInputName().equals(inputName)) {
throw new IllegalArgumentException("Input step with id " + nodeId + " does not accept input with name "
+ inputName + " (ACCEPTED NAME is " + step.getInputName() + ")");
}
step.setFieldValue(input);
if (eo.getInputSteps().stream().allMatch(s -> step.isReady()))
eo.getContext().setStatus(Status.READY);
else
eo.getContext().setStatus(Status.INITIALIZING);
} else if (eo.getContext().getStatus() == Status.WAITING) {
ExecutionStep step = eo.getExecutionSteps().stream().filter(s -> s.getId().equals(nodeId))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found"));
if (!(step instanceof InteractiveStep))
throw new IllegalArgumentException("step with id " + nodeId + " is not an interactive step");
InteractiveStep interactiveStep = (InteractiveStep) step;
if (!eo.getContext().getWaitingSteps().contains(interactiveStep.getId()))
throw new IllegalArgumentException("step with id " + nodeId + " is not in waiting state");
interactiveStep.setUserInput(input);
} else
throw new IllegalStateException("Execution with id " + executionId
+ " is not in initialization status (CURRENT STATUS is " + eo.getContext().getStatus() + ")");
}
InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found"));
if (!step.getInputName().equals(inputName)) {
throw new IllegalArgumentException("Input step with id " + nodeId + " does not accept input with name "
+ inputName + " (ACCEPTED NAME is " + step.getInputName() + ")");
}
step.setFieldValue(input);
if (eo.getInputSteps().stream().allMatch(s -> step.isReady()))
eo.getContext().setStatus(Status.READY);
else
eo.getContext().setStatus(Status.INITIALIZING);
return eo;
return eo;
}
public ExecutionObject startExecution(String id) {
ExecutionObject eo = getExecution(id);
if (eo.getContext().getStatus() != Status.READY) {
throw new IllegalStateException("Execution with id " + id + " is not in READY status (CURRENT STATUS is "
+ eo.getContext().getStatus() + ")");
}
eo.getInputSteps().forEach(InputStep::start);
eo.getContext().setStatus(Status.RUNNING);
eo.getInputSteps().forEach(InputStep::start);
return eo;
}

View File

@ -10,12 +10,11 @@ import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.entities.FlowEntity;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.auth.LoginEntity;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.repositories.AuthRepository;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import jakarta.annotation.PostConstruct;
@Component
public class ImportComponent {

View File

@ -11,6 +11,11 @@ import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.dto.Connection;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.OutputNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
@ -20,13 +25,11 @@ import it.cnr.isti.workflow.manager.dto.IOModel;
import it.cnr.isti.workflow.manager.model.flows.errors.FlowError;
import it.cnr.isti.workflow.manager.model.steps.ExecutionStep;
import it.cnr.isti.workflow.manager.model.steps.InputStep;
import it.cnr.isti.workflow.manager.model.steps.InteractiveStep;
import it.cnr.isti.workflow.manager.model.steps.OutputStep;
import it.cnr.isti.workflow.manager.model.steps.Step;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.model.steps.UserDefinedStep;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition;
import static it.cnr.isti.workflow.manager.model.flows.errors.FlowError.error;
import static it.cnr.isti.workflow.manager.model.flows.errors.ErrorType.*;
@ -62,20 +65,35 @@ public class TransformerService {
Step step = switch (node.getNodeDefinition()) {
case InputNodeDefinition id -> {
String inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY);
String inputName;
if (node.getParameters()==null)
inputName = (String) id.getRuntimeParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY).getDefaultValue();
else inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY);
IOType type = IOType.fromString(node.getOutputs().getFirst().getType());
outputs.put(node.getOutputs().getFirst().getKey(), new IDPair(node.getKey(), inputName));
yield InputStep.builder().id(node.getKey()).position(node.getPosition()).inputName(inputName).type(type).build();
}
case UserNodeDefinition ud -> {
ExecutionStep execStep = getExecutionStep(node, ud, execObject.getContext());
UserDefinedStep execStep = getUserDefinedStep(node, ud, execObject.getContext());
node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName())));
node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName())));
yield execStep;
}
case HumanInteractionNodeDefinition uid -> {
InteractiveStep execStep = getInteractiveStep(node, uid, execObject.getContext());
node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName())));
node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName())));
yield execStep;
}
case OutputNodeDefinition od -> {
String outputName;
if (node.getParameters()==null)
outputName = (String) od.getRuntimeParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY).getDefaultValue();
else outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY);
IOType type = IOType.fromString(node.getInputs().getFirst().getType());
String outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY);
inputs.put(node.getInputs().getFirst().getKey(), new IDPair(node.getKey(), outputName));
yield OutputStep.builder().id(node.getKey()).position(node.getPosition()).type(type).outputName(outputName).build();
}
@ -97,8 +115,7 @@ public class TransformerService {
Step outputNode = steps.get(outputPair.nodeId);
InputReadyListener inputNode = (InputReadyListener) steps.get(inputPair.nodeId);
outputNode.getNextSteps().add(inputNode);
inputNode.getParentOutputsToInputMapping().put(outputPair.name, inputPair.name);
inputNode.addParentOutputsToInputMapping(outputPair.name, inputPair.name);
outputNode.addInputMapping(outputPair.name, new FieldKey(inputPair.nodeId, inputPair.name));
}
@ -120,11 +137,11 @@ public class TransformerService {
return execObject;
}
private ExecutionStep getExecutionStep(Node node, UserNodeDefinition nodeDefinition, ExecutionContext context) {
private UserDefinedStep getUserDefinedStep(Node node, UserNodeDefinition nodeDefinition, ExecutionContext context) {
if (!executors.containsKey(nodeDefinition.getExecutor()))
throw new IllegalArgumentException("Executor not found for node type: " + node.getType());
ExecutionStep step = ExecutionStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition())
UserDefinedStep step = UserDefinedStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition())
.inputTypes(nodeDefinition.getInputs())
.outputTypes(nodeDefinition.getOutputs())
.executor(executors.get(nodeDefinition.getExecutor()))
@ -133,6 +150,21 @@ public class TransformerService {
return step;
}
private InteractiveStep getInteractiveStep(Node node, HumanInteractionNodeDefinition nodeDefinition, ExecutionContext context) {
IOType inputType = IOType.fromString(node.getInputs().getFirst().getType());
IOType outputType = IOType.fromString(node.getOutputs().getFirst().getType());
String actionDescription = (String) node.getParameters().get(HumanInteractionNodeDefinition.PARAMETER_NAME);
InteractiveStep step = InteractiveStep.builder().id(node.getKey()).position(node.getPosition())
.inputType(inputType)
.outputType(outputType)
.context(context)
.interactionDescription(actionDescription)
.build();
return step;
}
public List<FlowError> isTrasformableFlow(Flow flow) {
List<FlowError> errors = new ArrayList<>();
if (flow.getNodes().isEmpty())

View File

@ -1,6 +1,5 @@
package it.cnr.isti.workflow.manager;
import org.aspectj.weaver.ast.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.context.annotation.Bean;

View File

@ -9,9 +9,13 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.Import;
import org.springframework.test.context.TestPropertySource;
import org.springframework.util.ResourceUtils;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.HumanInteractionNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.services.ExecutionService;
@ -58,4 +62,35 @@ public class ExecutionServiceTest {
});
}
@Test()
void testInteractionExecutionRunningSuccess() throws Exception{
ObjectMapper objectMapper = new ObjectMapper();
Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:interactionFlow.json"), Flow.class);
ExecutionObject execution = executionService.createExecution(flow);
ObjectMapper mapper = new ObjectMapper();
log.info("Execution created: {}", mapper.writeValueAsString(execution));
execution = executionService.prepareInput(execution.getId(), "59d68aed-6de8-4305-9a3c-249069474d2e", InputNodeDefinition.DEFAULT_VALUE, "ciao como stoi");
execution = executionService.startExecution(execution.getId());
while (execution.getContext().getStatus().isRunningState())
try {
log.info("Execution status: {}", execution.getContext().getStatus());
Thread.sleep(1000);
} catch (InterruptedException e) { }
assert execution.getContext().getStatus() == Status.WAITING;
String waitingStepId = execution.getContext().getWaitingSteps().getFirst();
log.info("Waiting step id: {}", waitingStepId);
executionService.prepareInput(execution.getId(), waitingStepId, HumanInteractionNodeDefinition.INPUT_NAME, "ciao come stai?");
assert execution.getContext().getResult().size() == 1;
execution.getContext().getResult().forEach((k, v) -> {
log.info("Output: {} = {}", k, v);
});
}
}

View File

@ -88,10 +88,19 @@ public class FlowTest {
void testInvalidTrasformable() throws Exception {
ObjectMapper objectMapper = new ObjectMapper();
Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:invalidFlow.json"), Flow.class);
List<FlowError> errors = transformerService.isTrasformableFlow(flow);
assertNotNull(errors);
assertTrue(!errors.isEmpty());
}
@Test
void testInteractionFlowTrasformable() throws Exception {
ObjectMapper objectMapper = new ObjectMapper();
Flow flow = objectMapper.readValue(ResourceUtils.getFile("classpath:interactionFlow.json"), Flow.class);
List<FlowError> errors = transformerService.isTrasformableFlow(flow);
assertNotNull(errors);
log.debug("Errors: {}", errors);
assertTrue(errors.isEmpty());
}
}

View File

@ -8,10 +8,10 @@ import it.cnr.isti.workflow.manager.dto.Connection;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.dto.IOModel;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.Translator;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
@ -93,8 +93,6 @@ public class TestFlowCreator {
new Connection(UUID.randomUUID().toString(), execOutputKey, outputInputKey)
));
flow.getNodes().forEach(Node::resolveNodeDefinition);
return flow;
}
}

View File

@ -6,12 +6,12 @@ import org.springframework.test.context.TestPropertySource;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.SystemNodeDefinition;
import it.cnr.isti.workflow.manager.entities.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.ParameterDefinition;
import it.cnr.isti.workflow.manager.model.types.ParameterType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import org.springframework.beans.factory.annotation.Autowired;
@ -83,7 +83,7 @@ class NodeDefRepoTest {
private NodeDefinitionRepository nodeRepository;
@Autowired
private InputDefinitionRepository inputNodeRepository;
private SystemNodeDefinitionRepository systemNodeRepository;
@Test
@ -91,7 +91,7 @@ class NodeDefRepoTest {
NodeDefinition node = UserNodeDefinition.builder().name("DefectDetection")
.executor("GENERIC-AI").category("DefectDetection").createdBy("lucio.lelii").input("file", IOType.CSV)
.fixedParameters(Map.of("LLM",
ParameterDefinition.builder().name("param").label("param").type(ParameterType.Text).build()))
ParameterDefinition.builder().name("param").label("param").defaultValue("google-gemini").type(ParameterType.Text).build()))
.output("file", IOType.CSV).build();
nodeRepository.save(node);
NodeDefinition foundNode = nodeRepository.findById(node.getName()).orElse(null);
@ -113,7 +113,7 @@ class NodeDefRepoTest {
@Test
void getSystemNodeDefinitions() {
Iterable<InputNodeDefinition> nodes = inputNodeRepository.findAll();
Iterable<SystemNodeDefinition> nodes = systemNodeRepository.findAll();
assertNotNull(nodes);
long count = StreamSupport.stream(nodes.spliterator(), false).count();
//assertTrue(count == IOType.values().length);
@ -125,7 +125,7 @@ class NodeDefRepoTest {
NodeDefinition node = UserNodeDefinition.builder().name("DefectDetection")
.executor("GENERIC-AI").category("DefectDetection").createdBy("lucio.lelii").input("file", IOType.CSV)
.fixedParameters(Map.of("LLM",
ParameterDefinition.builder().name("param").label("param").type(ParameterType.Text).build()))
ParameterDefinition.builder().name("param").label("param").defaultValue("google-gemini").type(ParameterType.Text).build()))
.output("file", IOType.CSV).build();
nodeRepository.save(node);
NodeDefinition foundNode = nodeRepository.findById(node.getName()).orElse(null);

View File

@ -0,0 +1,168 @@
{
"createdBy": "lucio.lelii",
"name": "interaction flow",
"description": null,
"nodes": [
{
"key": "0414f102-fa57-47a5-994d-575bee0443fa",
"name": "newNode",
"createdBy": null,
"outputs": [
{
"key": "8dda93d6-7bc6-485e-8375-7c56ca676979",
"type": "TEXT",
"name": "output"
}
],
"inputs": [
{
"key": "57a64d4c-bbbf-4c56-a0d2-7a2b6de2da09",
"type": "TEXT",
"name": "input"
}
],
"color": "#A8E6CF",
"position": {
"x": 325,
"y": 145
},
"parameters": {
"action_description": "correct the phrase"
},
"description": null,
"nodeDefinition": {
"type": "system",
"name": "Human Interaction",
"createdBy": "system",
"category": "humanInteraction",
"color": "black",
"fixedParameters": {},
"runtimeParameters": {
"action_description": {
"name": "action_description",
"label": "Human Readable Action Description",
"description": "Description of the user action to perform",
"type": "Text",
"required": true,
"defaultValue": null,
"validations": [],
"specificAttributes": null
}
},
"inputs": {
"input": "TEXT"
},
"outputs": {
"output": "TEXT"
},
"simulable": false
},
"type": "Human Interaction"
},
{
"key": "59d68aed-6de8-4305-9a3c-249069474d2e",
"name": "newNode",
"createdBy": null,
"outputs": [
{
"key": "05d80c69-0c16-454a-ac81-ae8fa1278667",
"type": "TEXT",
"name": "value"
}
],
"inputs": [],
"color": "#A8E6CF",
"position": {
"x": 99,
"y": 214
},
"parameters": null,
"description": null,
"nodeDefinition": {
"type": "system",
"name": "Text Input",
"createdBy": null,
"category": "input",
"color": "green",
"fixedParameters": {},
"runtimeParameters": {
"name": {
"name": "name",
"label": "Input Name",
"description": "Input name",
"type": "String",
"required": false,
"defaultValue": "output",
"validations": [],
"specificAttributes": null
}
},
"inputs": null,
"outputs": {
"value": "TEXT"
},
"simulable": false
},
"type": "Text Input"
},
{
"key": "f9095513-fa24-4238-9b08-2a2b959e3978",
"name": "newNode",
"createdBy": null,
"outputs": [],
"inputs": [
{
"key": "8daefe93-5998-4362-b24a-5087c1b2f859",
"type": "TEXT",
"name": "value"
}
],
"color": "#A8E6CF",
"position": {
"x": 697,
"y": 200
},
"parameters": null,
"description": null,
"nodeDefinition": {
"type": "system",
"name": "Text Output",
"createdBy": null,
"category": "output",
"color": "red",
"fixedParameters": {},
"runtimeParameters": {
"name": {
"name": "name",
"label": "Output Name",
"description": "Output name",
"type": "String",
"required": false,
"defaultValue": "output",
"validations": [],
"specificAttributes": null
}
},
"inputs": {
"value": "TEXT"
},
"outputs": null,
"simulable": false
},
"type": "Text Output"
}
],
"connections": [
{
"key": "e8474793-a8d9-4fd2-9da6-1a7c5b8657d2",
"from": "05d80c69-0c16-454a-ac81-ae8fa1278667",
"to": "57a64d4c-bbbf-4c56-a0d2-7a2b6de2da09"
},
{
"key": "7f21e4aa-d554-4fdf-ba95-c331a48ccf92",
"from": "8dda93d6-7bc6-485e-8375-7c56ca676979",
"to": "8daefe93-5998-4362-b24a-5087c1b2f859"
}
],
"public": false
}