refactored a lot

This commit is contained in:
Lucio Lelii 2025-07-24 15:33:39 +02:00
parent c01271e98b
commit 4200a2bc3d
36 changed files with 498 additions and 234 deletions

View File

@ -7,15 +7,15 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Import;
import it.cnr.isti.workflow.manager.executors.ai.models.AIModel;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.OutputNodeDefinition;
import it.cnr.isti.workflow.manager.repositories.AuthRepository;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import it.cnr.isti.workflow.manager.services.ImportComponent;
@SpringBootApplication
public class WorkflowManagerApplication {
@ -23,13 +23,16 @@ public class WorkflowManagerApplication {
@Autowired
List<AIModel> aiModels;
@Autowired
ImportComponent importComponent;
public static void main(String[] args) {
SpringApplication.run(WorkflowManagerApplication.class, args);
}
@Bean
@ConditionalOnProperty(prefix = "app", name = "db.init.enabled", havingValue = "true")
CommandLineRunner init(NodeDefinitionRepository repository, FlowRepository flowRepository, AuthRepository authRepository) {
CommandLineRunner init(NodeDefinitionRepository repository) {
return args -> {
NodeDefinition textInput = InputNodeDefinition.builder().name("Text Input")
@ -41,6 +44,8 @@ public class WorkflowManagerApplication {
.build();
repository.save(textInput);
repository.save(textOutput);
importComponent.start();
};
}

View File

@ -18,10 +18,10 @@ import org.springframework.web.multipart.MultipartFile;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.exceptions.WebException;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;
import it.cnr.isti.workflow.manager.services.ExecutionService;
@ -178,11 +178,11 @@ public class ExecutionController {
* @return An {@link ExecutionObject} representing the updated execution state
* after preparing the input.
*/
@PutMapping(path = "{id}/input/{inputName}")
@PutMapping(path = "{id}/node/{nodeId}/input/{inputName}/text", consumes = "text/plain")
@Operation(summary = "Prepares string inputs", description = "Prepares the input for an execution by associating a given input string with a specific input name and execution ID.")
public ExecutionObject prepareStringInputs(@RequestBody String input, @PathVariable String inputName,
public ExecutionObject prepareStringInputs(@RequestBody String input, @PathVariable String nodeId, @PathVariable String inputName,
@PathVariable String id) {
return executionService.prepareInput(id, inputName, input);
return executionService.prepareInput(id, nodeId, inputName, input);
}
@ -197,14 +197,14 @@ public class ExecutionController {
* preparing the input.
* @throws WebServerException If an error occurs while creating or transferring the file.
*/
@PutMapping(path = "{id}/input/{inputName}", consumes = "multipart/form-data")
@PutMapping(path = "{id}/node/{nodeId}/input/{inputName}/file", consumes = "multipart/form-data")
@Operation(summary = "Prepares file inputs", description = "Prepares file inputs for a specific execution by uploading a file and associating it with the given input name and execution ID.")
public ExecutionObject prepareFileInputs(@RequestParam("file") MultipartFile file, @PathVariable String inputName,
public ExecutionObject prepareFileInputs(@RequestParam("file") MultipartFile file, @PathVariable String nodeId, @PathVariable String inputName,
@PathVariable String id) {
try {
File myFile = File.createTempFile(inputName, file.getOriginalFilename());
file.transferTo(myFile);
return executionService.prepareInput(id, inputName, myFile);
return executionService.prepareInput(id, nodeId, inputName, myFile);
} catch (IOException e) {
throw new WebServerException("Error while creating file", e);
}

View File

@ -18,7 +18,8 @@ import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.exceptions.WebException;
import it.cnr.isti.workflow.manager.model.auth.LoginEntity;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.entities.FlowEntity;
import it.cnr.isti.workflow.manager.model.flows.errors.FlowError;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;
import it.cnr.isti.workflow.manager.services.TransformerService;
@ -27,7 +28,6 @@ import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
@ -90,7 +90,7 @@ public class FlowsController {
if (flow.getId() != null)
throw new WebException(HttpStatus.BAD_REQUEST, "on creation flow id must be null");
flow.setCreatedBy(userDetails.getUsername());
Flow savedFlow = flowRepository.save(flow);
Flow savedFlow = flowRepository.save((FlowEntity)flow);
logger.info("flow created {}", flow.getId());
return savedFlow;
@ -118,7 +118,7 @@ public class FlowsController {
if (!savedFlow.getCreatedBy().equals(userDetails.getUsername()))
throw new WebException(HttpStatus.FORBIDDEN, "Flow not created by user");
savedFlow = flowRepository.save(flow);
savedFlow = flowRepository.save(FlowEntity.from(flow));
logger.info("Flow saved: {} ",savedFlow.getId());
return savedFlow;
}

View File

@ -1,4 +1,4 @@
package it.cnr.isti.workflow.manager.model.flows;
package it.cnr.isti.workflow.manager.dto;
import lombok.AllArgsConstructor;
import lombok.Builder;

View File

@ -0,0 +1,40 @@
package it.cnr.isti.workflow.manager.dto;
import java.util.List;
import com.fasterxml.jackson.annotation.JsonAlias;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Builder.Default;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.NonNull;
import lombok.ToString;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
@ToString
public class Flow {
String id;
@NonNull
String createdBy;
@Default
@JsonAlias("public")
boolean isPublic =false;
@NonNull
String name;
String description;
@Builder.Default
List<Node> nodes = List.of();
@Builder.Default
List<Connection> connections = List.of();
}

View File

@ -1,4 +1,4 @@
package it.cnr.isti.workflow.manager.model.flows;
package it.cnr.isti.workflow.manager.dto;
import lombok.AllArgsConstructor;
import lombok.Data;

View File

@ -0,0 +1,58 @@
package it.cnr.isti.workflow.manager.dto;
import java.util.List;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.services.RepositoryHolder;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.Singular;
import lombok.ToString;
@Data
@Builder
@AllArgsConstructor
@NoArgsConstructor
@ToString
public class Node {
String key;
String name;
String createdBy;
@Singular
List<IOModel> outputs;
@Singular
List<IOModel> inputs;
String color;
Point position;
Map<String, Object> parameters;
String description;
@JsonProperty(access = JsonProperty.Access.READ_ONLY)
NodeDefinition nodeDefinition;
String type;
public void resolveNodeDefinition() {
if (type != null) {
this.nodeDefinition = RepositoryHolder.getNodeDefinition()
.findById(type)
.orElse(null);
}
}
}

View File

@ -0,0 +1,13 @@
package it.cnr.isti.workflow.manager.dto;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@NoArgsConstructor
@AllArgsConstructor
@Data
public class Point {
int x;
int y;
}

View File

@ -0,0 +1,128 @@
package it.cnr.isti.workflow.manager.entities;
import java.util.List;
import org.hibernate.annotations.UuidGenerator;
import it.cnr.isti.workflow.manager.dto.Connection;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.repositories.converters.ConnectionListConverter;
import it.cnr.isti.workflow.manager.repositories.converters.NodeListConverter;
import jakarta.persistence.Column;
import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
import jakarta.persistence.Id;
import jakarta.persistence.PostLoad;
import lombok.EqualsAndHashCode;
@Entity
@EqualsAndHashCode(callSuper = true)
public class FlowEntity extends Flow {
public FlowEntity() {
super();
}
public FlowEntity(Flow flow) {
super(
flow.getId(),
flow.getCreatedBy(),
flow.isPublic(),
flow.getName(),
flow.getDescription(),
flow.getNodes(),
flow.getConnections()
);
}
public static FlowEntity from(Flow flow) {
return new FlowEntity(flow);
}
@Id
@UuidGenerator
@Override
public String getId() {
return super.getId();
}
@Override
public void setId(String id) {
super.setId(id);
}
@Column(nullable = false)
@Override
public String getCreatedBy() {
return super.getCreatedBy();
}
@Override
public void setCreatedBy(String createdBy) {
super.setCreatedBy(createdBy);
}
@Column(name = "public", nullable = false)
@Override
public boolean isPublic() {
return super.isPublic();
}
@Override
public void setPublic(boolean isPublic) {
super.setPublic(isPublic);
}
@Column(nullable = false)
@Override
public String getName() {
return super.getName();
}
@Override
public void setName(String name) {
super.setName(name);
}
@Column
@Override
public String getDescription() {
return super.getDescription();
}
@Override
public void setDescription(String description) {
super.setDescription(description);
}
@Convert(converter = NodeListConverter.class)
@Column(columnDefinition = "TEXT")
@Override
public List<Node> getNodes() {
return super.getNodes();
}
@Override
public void setNodes(List<Node> nodes) {
super.setNodes(nodes);
}
@Convert(converter = ConnectionListConverter.class)
@Column(columnDefinition = "TEXT")
@Override
public List<Connection> getConnections() {
return super.getConnections();
}
@Override
public void setConnections(List<Connection> connections) {
super.setConnections(connections);
}
@PostLoad
public void postLoad() {
this.getNodes().forEach(Node::resolveNodeDefinition);
}
}

View File

@ -7,9 +7,11 @@ import java.util.UUID;
import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.FieldKey;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.model.steps.ExecutionStep;
import it.cnr.isti.workflow.manager.model.steps.InputStep;
import it.cnr.isti.workflow.manager.model.steps.OutputStep;
import it.cnr.isti.workflow.manager.model.steps.ResultListener;
@ -30,30 +32,28 @@ public class ExecutionObject implements ResultListener {
@Setter
String name;
Map<String , IOType> exepectedOutputsDescription;
Map<FieldKey , IOType> exepectedOutputsDescription;
Map<String , IOType> exepectedInputsDescription;
Map<FieldKey , IOType> exepectedInputsDescription;
Flow flow;
@JsonIgnore
List<InputStep> inputSteps;
@JsonIgnore
List<OutputStep> outputSteps;
List<ExecutionStep> executionSteps;
public ExecutionObject(Flow flow) {
this.flow = flow;
this.name = flow.getName()+" - "+this.creationTime;
public ExecutionObject(String executionName) {
this.name = executionName;
}
public void setOutputSteps(List<OutputStep> outputSteps) {
this.outputSteps = outputSteps;
this.exepectedOutputsDescription = new HashMap<>();
for (OutputStep outputStep : outputSteps) {
this.exepectedOutputsDescription.put(outputStep.getOutputName(), outputStep.getType());
this.exepectedOutputsDescription.put(new FieldKey(outputStep.getId(), outputStep.getOutputName()), outputStep.getOutputType());
outputStep.registerResultListener(this);
}
}
@ -61,8 +61,12 @@ public class ExecutionObject implements ResultListener {
public void setInputSteps(List<InputStep> inputSteps) {
this.inputSteps = inputSteps;
this.exepectedInputsDescription = new HashMap<>();
for (InputStep inputStep : inputSteps)
this.exepectedInputsDescription.put(inputStep.getInputName(), inputStep.getType());
for (InputStep inputStep : inputSteps)
this.exepectedInputsDescription.put(new FieldKey(inputStep.getId(), inputStep.getInputName()), inputStep.getInputType());
}
public void setExecutionSteps(List<ExecutionStep> executionSteps) {
this.executionSteps = executionSteps;
}
@Override
@ -70,7 +74,7 @@ public class ExecutionObject implements ResultListener {
//TODO: devo discriminare in base al nodeId altrimenti se ho più output step con lo stesso inputName non so a quale output step riferirmi
this.context.addResult(resultKey, value);
this.context.addResult(nodeKey, resultKey, value);
if (this.exepectedOutputsDescription.keySet().containsAll(this.context.getExecutionResult().keySet())) {
//significa che l'esecuzione è finita
this.context.setStatus(Status.SUCCESS);

View File

@ -47,7 +47,7 @@ public class GeminiService {
status -> status.is5xxServerError(),
clientResponse -> clientResponse.bodyToMono(String.class)
.defaultIfEmpty("Error: server without body")
.flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body))))
.flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body))))
.bodyToMono(String.class)
.timeout(Duration.ofMinutes(2))
.retryWhen(

View File

@ -67,7 +67,7 @@ public class OllamaService {
status -> status.is5xxServerError(),
clientResponse -> clientResponse.bodyToMono(String.class)
.defaultIfEmpty("Error: server without body")
.flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body))))
.flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body))))
.bodyToMono(GenerateResponse.class) // deserialize JSON in oggetto Java
.timeout(Duration.ofMinutes(2));
@ -86,7 +86,7 @@ public class OllamaService {
status -> status.is5xxServerError(),
clientResponse -> clientResponse.bodyToMono(String.class)
.defaultIfEmpty("Error: server without body")
.flatMap(body -> Mono.error(new RuntimeException("Errore 5xx: " + body))))
.flatMap(body -> Mono.error(new RuntimeException("Error 5xx: " + body))))
.bodyToMono(ModelResponse.class) // deserialize JSON in oggetto Java
.timeout(Duration.ofMinutes(1))
.map(modelResponse -> modelResponse.getModels().stream()

View File

@ -5,8 +5,6 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import lombok.AccessLevel;
import lombok.Getter;
import lombok.NoArgsConstructor;
@ -47,14 +45,14 @@ public class ExecutionContext {
Map<String, Object> inputs = new HashMap<>();
@Getter()
private Map<String, Object> result = new HashMap<>();
private Map<FieldKey, Object> result = new HashMap<>();
Long startTime = null;
Long endTime = null;
List<String> stepsUnderExecution = new ArrayList<>();
List<String> errors = new ArrayList<>();;
List<String> warnings = new ArrayList<>();;
Map<String, String> errors = new HashMap<>();
List<String> warnings = new ArrayList<>();
Map<String, Map<String, Object>> nodeResult = new HashMap<>();
@ -74,15 +72,19 @@ public class ExecutionContext {
this.nodeResult.put(nodeId, result);
}
public void addResult(String key, Object value) {
this.result.put(key, value);
public void addResult(String nodeId, String key, Object value) {
this.result.put(new FieldKey(nodeId, key), value);
}
public void addInput(String key, Object value) {
this.inputs.put(key, value);
}
public Map<String, Object> getExecutionResult() {
public void addError(String nodeId, String error) {
this.errors.put(nodeId, error);
}
public Map<FieldKey, Object> getExecutionResult() {
return Collections.unmodifiableMap(this.result);
}
}

View File

@ -0,0 +1,8 @@
package it.cnr.isti.workflow.manager.model;
public record FieldKey(String nodeId, String fieldId) {
@Override
public final String toString() {
return nodeId + ":" + fieldId;
}
}

View File

@ -1,63 +0,0 @@
package it.cnr.isti.workflow.manager.model.flows;
import java.util.List;
import org.hibernate.annotations.UuidGenerator;
import com.fasterxml.jackson.annotation.JsonAlias;
import it.cnr.isti.workflow.manager.repositories.converters.ConnectionListConverter;
import it.cnr.isti.workflow.manager.repositories.converters.NodeListConverter;
import jakarta.persistence.Column;
import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
import jakarta.persistence.Id;
import jakarta.validation.constraints.NotBlank;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Builder.Default;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.NonNull;
import lombok.ToString;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
@ToString
@Entity
public class Flow {
@Id
@UuidGenerator
String id;
@NotBlank
@NonNull
@Column(name = "created_by")
String createdBy;
@Default
@JsonAlias("public")
boolean isPublic =false;
@NotBlank
@NonNull
String name;
String description;
@Builder.Default
@Convert(converter = NodeListConverter.class)
@Column(columnDefinition = "TEXT")
List<Node> nodes = List.of();
@Builder.Default
@Convert(converter = ConnectionListConverter.class)
@Column(columnDefinition = "TEXT")
List<Connection> connections = List.of();
}

View File

@ -1,50 +0,0 @@
package it.cnr.isti.workflow.manager.model.flows;
import java.util.List;
import java.util.Map;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.Singular;
import lombok.ToString;
@Data
@Builder
@AllArgsConstructor
@NoArgsConstructor
@ToString
public class Node {
String key;
String name;
String createdBy;
@Singular
List<IOModel> outputs;
@Singular
List<IOModel> inputs;
String color;
Point position;
Map<String, Object> parameters;
String description;
String type;
}
@NoArgsConstructor
@AllArgsConstructor
@Data
class Point {
int x;
int y;
}

View File

@ -6,17 +6,17 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import it.cnr.isti.workflow.manager.dto.Point;
import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.exceptions.ExecutionException;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.Translator;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.EqualsAndHashCode;
@ -24,40 +24,53 @@ import lombok.NonNull;
import lombok.ToString;
@Data
@AllArgsConstructor
@ToString
@EqualsAndHashCode(callSuper = true)
public class ExecutionStep extends Step implements InputReadyListener, ExecutableStep {
private static final Logger log = org.slf4j.LoggerFactory.getLogger(ExecutionStep.class);
@NonNull
String name;
Map<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(ExecutionContext context, String id, Map<String, Object> runtimeParameters,
public ExecutionStep(@NonNull String name, @NonNull ExecutionContext context, @NonNull String id, Map<String, Object> runtimeParameters,
Map<String, Object> fixedParameters,
UserNodeDefinition nodeDefinition,
Executor executor) {
this.id = id;
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
@ -98,12 +111,12 @@ public class ExecutionStep extends Step implements InputReadyListener, Executabl
} catch (ExecutionException e) {
log.error("error during execution of step {}: {}", id, e.getMessage());
this.context.getErrors().add(String.format("[%S] %s", this.id, e.getMessage()));
this.context.addError(this.id, e.getMessage());
this.context.setStatus(Status.ERROR);
return;
} finally {
this.context.getStepsUnderExecution().remove(this.id);
log.info("finished executing of step {}", this.nodeDefinition.getName());
log.info("finished executing step {}", this.nodeDefinition.getName());
}
}).start();

View File

@ -1,5 +1,7 @@
package it.cnr.isti.workflow.manager.model.steps;
import it.cnr.isti.workflow.manager.dto.Point;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import lombok.Builder;
@ -11,17 +13,17 @@ import lombok.NonNull;
@EqualsAndHashCode(callSuper = true)
public class InputStep extends Step implements ExecutableStep {
IOType type;
IOType inputType;
Object fieldValue;
String inputName;
@Builder
InputStep(@NonNull String id,@NonNull String inputName, @NonNull IOType type) {
this.id = id;
InputStep(@NonNull String id,@NonNull String inputName, @NonNull IOType type, Point position) {
super(id, position);
this.inputName = inputName;
this.type = type;
this.inputType = type;
}
public void setFieldValue(Object fieldValue) {
@ -35,7 +37,7 @@ public class InputStep extends Step implements ExecutableStep {
@Override
public void start() {
this.getNextSteps().forEach(step -> {
step.inputReady(InputNodeDefinition.OUTPUT_NAME, this.fieldValue);
step.inputReady(this.inputName, this.fieldValue);
});
}

View File

@ -5,6 +5,11 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
import io.micrometer.common.lang.NonNull;
import it.cnr.isti.workflow.manager.dto.Point;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.types.IOType;
import lombok.Builder;
@ -15,21 +20,20 @@ import lombok.EqualsAndHashCode;
@EqualsAndHashCode(callSuper = true)
public class OutputStep extends Step implements InputReadyListener, ResultProvider {
private String id;
private IOType type;
private IOType outputType;
private String outputName;
@JsonIgnore
List<ResultListener> resultListeners = new ArrayList<>();
Map<String, String> parentOutputsToInputMapping = new HashMap<>();
@Builder
OutputStep(String id, IOType type, String outputName) {
OutputStep(@NonNull String id, @NonNull IOType type, @NonNull String outputName, Point position) {
super(id, position);
this.outputName = outputName;
this.id = id;
this.type = type;
this.outputType = type;
}
@Override

View File

@ -1,16 +1,50 @@
package it.cnr.isti.workflow.manager.model.steps;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import it.cnr.isti.workflow.manager.dto.Point;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import io.micrometer.common.lang.NonNull;
import it.cnr.isti.workflow.manager.model.FieldKey;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import lombok.Data;
@Data
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "type")
@JsonSubTypes({
@JsonSubTypes.Type(value = InputStep.class, name = "input"),
@JsonSubTypes.Type(value = OutputStep.class, name = "output"),
@JsonSubTypes.Type(value = ExecutionStep.class, name = "execution")
})
public abstract class Step {
@NonNull
String id;
//output name to external input for connections
Map<String, FieldKey> inputMapping = new HashMap<>();
Point position;
@JsonIgnore
List<InputReadyListener> nextSteps = new ArrayList<>();
boolean startAutomatically = true;
public Step(String id, Point position){
this.id = id;
this.position = position;
}
public void addInputMapping(String outputName, FieldKey fieldKey) {
this.inputMapping.put(outputName, fieldKey);
}
}

View File

@ -2,6 +2,8 @@ package it.cnr.isti.workflow.manager.model.types;
import java.io.IOException;
import java.util.List;
import org.apache.commons.lang3.Validate;
import org.json.JSONObject;
import com.fasterxml.jackson.core.JacksonException;
@ -38,8 +40,9 @@ public class ParameterDefinition {
@NonNull
private ParameterType type;
@Builder.Default
private boolean required = false;
protected boolean required = false;
private Object defaultValue;
@Singular
private List<Validation> validations;
@ -47,6 +50,22 @@ public class ParameterDefinition {
@JsonDeserialize(using = JsonObjectDeserializer.class)
@JsonSerialize(using = JsonObjectSerializer.class)
private JSONObject specificAttributes;
public static ParameterDefinitionBuilder builder() {
return new CustomParameterDefinitionBuilder();
}
// extend the generated builder class to add custom logic
private static class CustomParameterDefinitionBuilder extends ParameterDefinitionBuilder {
@Override
public ParameterDefinition build() {
if (!super.required && super.defaultValue == null) {
throw new IllegalArgumentException("When not required the parameter must have a default value.");
}
return super.build();
}
}
}
class JsonObjectDeserializer extends JsonDeserializer<JSONObject> {

View File

@ -5,6 +5,7 @@ public enum ParameterType {
Select,
DynamicMultiOptions,
Text,
String,
Boolean,
Number
}

View File

@ -30,9 +30,9 @@ public class InputNodeDefinition extends NodeDefinition {
this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder()
.name(INPUT_PARAMETER_KEY)
.label("Input Name")
.required(true)
.description("Input name")
.type(ParameterType.Text)
.defaultValue("output")
.type(ParameterType.String)
.build());
}

View File

@ -32,9 +32,9 @@ public class OutputNodeDefinition extends NodeDefinition {
this.runtimeParameters = Map.of(INPUT_PARAMETER_KEY, ParameterDefinition.builder()
.name(INPUT_PARAMETER_KEY)
.label("Output Name")
.required(true)
.defaultValue("output")
.description("Output name")
.type(ParameterType.Text)
.type(ParameterType.String)
.build());
}

View File

@ -7,11 +7,13 @@ import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
import org.springframework.stereotype.Repository;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.entities.FlowEntity;
@Repository
public interface FlowRepository extends JpaRepository<Flow, String> {
public interface FlowRepository extends JpaRepository<FlowEntity, String> {
@Query("SELECT f FROM Flow f WHERE f.createdBy = :createdBy OR f.isPublic = true")
@Query("SELECT f FROM FlowEntity f WHERE f.createdBy = :createdBy OR f.public = true")
List<Flow> findFlowsByCreatedByOrPublic(@Param("createdBy") String createdBy);
}

View File

@ -6,7 +6,7 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.model.flows.Connection;
import it.cnr.isti.workflow.manager.dto.Connection;
import jakarta.persistence.AttributeConverter;
import jakarta.persistence.Converter;

View File

@ -4,9 +4,8 @@ import java.util.List;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.model.flows.Node;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.services.ObjectMapperHolder;
import jakarta.persistence.AttributeConverter;
import jakarta.persistence.Converter;
@ -14,7 +13,7 @@ import jakarta.persistence.Converter;
@Converter
public class NodeListConverter implements AttributeConverter<List<Node>, String> {
@Override
public String convertToDatabaseColumn(List<Node> nodes) {

View File

@ -8,7 +8,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.model.steps.InputStep;
@Service
@ -19,7 +18,7 @@ public class ExecutionService {
@Autowired
TransformerService transformerService;
public ExecutionObject createExecution(Flow flow) {
public ExecutionObject createExecution(it.cnr.isti.workflow.manager.dto.Flow flow) {
ExecutionObject execObject = transformerService.transform(flow);
executions.put(execObject.getId(), execObject);
return execObject;
@ -44,17 +43,23 @@ public class ExecutionService {
executions.remove(id);
}
public ExecutionObject prepareInput(String id, String key, Object input) {
ExecutionObject eo = getExecution(id);
public ExecutionObject prepareInput(String executionId, String nodeId, String inputName, Object input) {
ExecutionObject eo = getExecution(executionId);
if (!eo.getContext().getStatus().isInitState()) {
throw new IllegalStateException("Execution with id " + id
throw new IllegalStateException("Execution with id " + executionId
+ " is not in initialization status (CURRENT STATUS is " + eo.getContext().getStatus() + ")");
}
InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(key))
InputStep step = eo.getInputSteps().stream().filter(s -> s.getId().equals(nodeId))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("Input step with id " + key + " not found"));
.orElseThrow(() -> new IllegalArgumentException("Input step with id " + nodeId + " not found"));
if (!step.getInputName().equals(inputName)) {
throw new IllegalArgumentException("Input step with id " + nodeId + " does not accept input with name "
+ inputName + " (ACCEPTED NAME is " + step.getInputName() + ")");
}
step.setFieldValue(input);
if (eo.getInputSteps().stream().allMatch(s -> step.isReady()))

View File

@ -8,9 +8,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.entities.FlowEntity;
import it.cnr.isti.workflow.manager.model.auth.LoginEntity;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
import it.cnr.isti.workflow.manager.repositories.AuthRepository;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;
@ -32,6 +32,9 @@ public class ImportComponent {
@Autowired
NodeDefinitionRepository nodeRepository;
@Autowired
RepositoryHolder repositoryHolder;
@Autowired
FlowRepository flowRepository;
@ -44,8 +47,7 @@ public class ImportComponent {
@Value("${app.import.path}")
String path;
@PostConstruct
void init() {
public void start() {
if (!enabled) {
System.out.println("ImportComponent is disabled");
return;
@ -85,7 +87,7 @@ public class ImportComponent {
logger.debug("Node {} not found, saving new node", node.getName());
nodeRepository.save(node);
}
);
);
} catch (Exception e) {
logger.error("Error loading nodes: {}", e.getMessage(),e);
}
@ -102,7 +104,7 @@ public class ImportComponent {
logger.info("Loading flows from file: {}", file.getAbsolutePath());
List<Flow> flows = ObjectMapperHolder.mapper.readValue(file, ObjectMapperHolder.mapper.getTypeFactory().constructCollectionType(List.class,Flow.class));
for (Flow flow: flows){
flowRepository.saveAndFlush(flow);
flowRepository.saveAndFlush(FlowEntity.from(flow));
}
} catch (Exception e) {
logger.error("Error loading flows: {}", e.getMessage(),e);

View File

@ -0,0 +1,26 @@
package it.cnr.isti.workflow.manager.services;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import jakarta.annotation.PostConstruct;
@Component
public class RepositoryHolder {
@Autowired
private NodeDefinitionRepository nodeDefinitionRepository;
private static NodeDefinitionRepository _nodeRepository;
@PostConstruct
public void init() {
System.out.println("Initializing RepositoryHolder with NodeDefinitionRepository");
RepositoryHolder._nodeRepository = nodeDefinitionRepository;
}
public static NodeDefinitionRepository getNodeDefinition() {
return _nodeRepository;
}
}

View File

@ -8,21 +8,21 @@ import java.util.Map;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import it.cnr.isti.workflow.manager.dto.Connection;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.executors.Executor;
import it.cnr.isti.workflow.manager.model.ExecutionContext;
import it.cnr.isti.workflow.manager.model.FieldKey;
import it.cnr.isti.workflow.manager.model.InputReadyListener;
import it.cnr.isti.workflow.manager.model.flows.Connection;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.model.flows.IOModel;
import it.cnr.isti.workflow.manager.model.flows.Node;
import it.cnr.isti.workflow.manager.dto.IOModel;
import it.cnr.isti.workflow.manager.model.flows.errors.FlowError;
import it.cnr.isti.workflow.manager.model.steps.ExecutionStep;
import it.cnr.isti.workflow.manager.model.steps.InputStep;
import it.cnr.isti.workflow.manager.model.steps.OutputStep;
import it.cnr.isti.workflow.manager.model.steps.Step;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.UserNodeDefinition;
import it.cnr.isti.workflow.manager.repositories.NodeDefinitionRepository;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.InputNodeDefinition;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
@ -40,15 +40,8 @@ public class TransformerService {
@Autowired
Map<String, Executor> executors;
@Autowired
NodeDefinitionRepository nodeDefinitionRepository;
private NodeDefinition getNodeDefinition(String nodeType) {
return nodeDefinitionRepository.findById(nodeType)
.orElseThrow(() -> new IllegalArgumentException("Node definition not found for type: " + nodeType));
}
public ExecutionObject transform(Flow flow) {
public ExecutionObject transform(it.cnr.isti.workflow.manager.dto.Flow flow) {
if (flow == null )
throw new IllegalArgumentException("Flow cannot be null");
@ -63,29 +56,33 @@ public class TransformerService {
Map<String, IDPair> inputs = new HashMap<>();
Map<String, IDPair> outputs = new HashMap<>();
ExecutionObject execObject = new ExecutionObject(flow);
ExecutionObject execObject = new ExecutionObject(flow.getName());
for (Node node : flow.getNodes()) {
Step step = switch (getNodeDefinition(node.getType())) {
Step step = switch (node.getNodeDefinition()) {
case InputNodeDefinition id -> {
String inputName = (String) node.getParameters().get(InputNodeDefinition.INPUT_PARAMETER_KEY);
IOType type = IOType.fromString(node.getOutputs().getFirst().getType());
yield InputStep.builder().id(node.getKey()).inputName(inputName).type(type).build();
outputs.put(node.getOutputs().getFirst().getKey(), new IDPair(node.getKey(), inputName));
yield InputStep.builder().id(node.getKey()).position(node.getPosition()).inputName(inputName).type(type).build();
}
case UserNodeDefinition ud -> {
ExecutionStep execStep = getExecutionStep(node, ud, execObject.getContext());
node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName())));
node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName())));
yield execStep;
}
case UserNodeDefinition ud -> getExecutionStep(node, ud, execObject.getContext());
case OutputNodeDefinition od -> {
IOType type = IOType.fromString(node.getInputs().getFirst().getType());
String outputName = (String) node.getParameters().get(OutputNodeDefinition.INPUT_PARAMETER_KEY);
yield OutputStep.builder().id(node.getKey()).type(type).outputName(outputName).build();
inputs.put(node.getInputs().getFirst().getKey(), new IDPair(node.getKey(), outputName));
yield OutputStep.builder().id(node.getKey()).position(node.getPosition()).type(type).outputName(outputName).build();
}
default -> throw new IllegalArgumentException("Node type " + node.getType() + " not supported");
};
steps.put(node.getKey(), step);
node.getOutputs().forEach(o -> outputs.put(o.getKey(), new IDPair(node.getKey(), o.getName())));
node.getInputs().forEach(i -> inputs.put(i.getKey(), new IDPair(node.getKey(), i.getName())));
}
for (Connection connection : flow.getConnections()) {
@ -101,6 +98,8 @@ public class TransformerService {
InputReadyListener inputNode = (InputReadyListener) steps.get(inputPair.nodeId);
outputNode.getNextSteps().add(inputNode);
inputNode.getParentOutputsToInputMapping().put(outputPair.name, inputPair.name);
outputNode.addInputMapping(outputPair.name, new FieldKey(inputPair.nodeId, inputPair.name));
}
execObject.setInputSteps(steps.values().stream()
@ -113,6 +112,11 @@ public class TransformerService {
.map(s -> (OutputStep) s)
.toList());
execObject.setExecutionSteps(steps.values().stream()
.filter(s -> s instanceof ExecutionStep)
.map(s -> (ExecutionStep) s)
.toList());
return execObject;
}
@ -120,7 +124,9 @@ public class TransformerService {
if (!executors.containsKey(nodeDefinition.getExecutor()))
throw new IllegalArgumentException("Executor not found for node type: " + node.getType());
ExecutionStep step = ExecutionStep.builder().id(node.getKey())
ExecutionStep step = ExecutionStep.builder().id(node.getKey()).name(nodeDefinition.getName()).position(node.getPosition())
.inputTypes(nodeDefinition.getInputs())
.outputTypes(nodeDefinition.getOutputs())
.executor(executors.get(nodeDefinition.getExecutor()))
.runtimeParameters(node.getParameters()).nodeDefinition(nodeDefinition).context(context).build();
@ -193,7 +199,7 @@ public class TransformerService {
// CHECK PARAMETERS
// TODO: devo controllare che i parametri siano del tipo giusto
NodeDefinition nodeDef = getNodeDefinition(n.getType());
NodeDefinition nodeDef = n.getNodeDefinition();
if (nodeDef.getRuntimeParameters() !=null && n.getParameters() != null) {
nodeDef.getRuntimeParameters().forEach((k, v) -> {

View File

@ -157,7 +157,7 @@
"y": -16
},
"parameters": {
"name": "interview\n"
"name": "interview"
},
"description": null,
"type": "Text Input"

View File

@ -10,6 +10,8 @@ import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.Import;
import org.springframework.test.context.TestPropertySource;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.services.ExecutionService;
@ -34,9 +36,11 @@ public class ExecutionServiceTest {
}
@Test()
void testExecutionRunningSuccess() {
void testExecutionRunningSuccess() throws Exception{
ExecutionObject execution = executionService.createExecution(testFlowCreator.create());
execution = executionService.prepareInput(execution.getId(), TestFlowCreator.INPUT_KEY, "ciao come stai");
ObjectMapper mapper = new ObjectMapper();
log.info("Execution created: {}", mapper.writeValueAsString(execution));
execution = executionService.prepareInput(execution.getId(), TestFlowCreator.INPUT_KEY, "input", "ciao come stai");
execution = executionService.startExecution(execution.getId());
while (execution.getContext().getStatus().isRunningState())

View File

@ -15,7 +15,7 @@ import org.springframework.util.ResourceUtils;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.executors.ExecutionObject;
import it.cnr.isti.workflow.manager.model.ExecutionContext.Status;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.model.flows.errors.FlowError;
import it.cnr.isti.workflow.manager.model.steps.ResultListener;
import it.cnr.isti.workflow.manager.repositories.FlowRepository;

View File

@ -4,10 +4,10 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
import it.cnr.isti.workflow.manager.model.flows.Connection;
import it.cnr.isti.workflow.manager.model.flows.Flow;
import it.cnr.isti.workflow.manager.model.flows.IOModel;
import it.cnr.isti.workflow.manager.model.flows.Node;
import it.cnr.isti.workflow.manager.dto.Connection;
import it.cnr.isti.workflow.manager.dto.Flow;
import it.cnr.isti.workflow.manager.dto.IOModel;
import it.cnr.isti.workflow.manager.dto.Node;
import it.cnr.isti.workflow.manager.model.types.IOType;
import it.cnr.isti.workflow.manager.model.types.Translator;
import it.cnr.isti.workflow.manager.model.types.nodes.definitions.NodeDefinition;
@ -26,7 +26,7 @@ public class TestFlowCreator {
}
public Flow create(){
Translator inputFieldtransaltor = Translator.builder().translation("translate the following phrase : ${{phrase}} in french, return only the translation").build();
Translator inputFieldtransaltor = Translator.builder().translation("translate the following phrase : ${{phrase}} in french, return only the translation").build();
NodeDefinition nFrenchTranslator = UserNodeDefinition.builder().name("FrenchTranslator")
.executor("GENERIC-AI").category("Translators").createdBy("lucio.lelii").input("phrase", IOType.TEXT)
@ -53,6 +53,7 @@ public class TestFlowCreator {
.input(new IOModel(execInputKey, IOType.TEXT.toString(), "phrase"))
.output(new IOModel(execOutputKey, IOType.TEXT.toString(), "translated"))
.parameters(parameter)
.nodeDefinition(nFrenchTranslator)
.createdBy("lucio.lelii")
.color("#FF0000")
.build();
@ -92,6 +93,8 @@ public class TestFlowCreator {
new Connection(UUID.randomUUID().toString(), execOutputKey, outputInputKey)
));
flow.getNodes().forEach(Node::resolveNodeDefinition);
return flow;
}
}

View File

@ -1,6 +1,5 @@
package it.cnr.isti.workflow.manager.executors;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;