diff --git a/pom.xml b/pom.xml
index 0ef8f38..a548e56 100644
--- a/pom.xml
+++ b/pom.xml
@@ -71,6 +71,13 @@
0.11.5
runtime
+
+ io.github.classgraph
+ classgraph
+ 4.8.159
+
+
+
com.kjetland
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java
index 0c187f0..8dd2236 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/Block.java
@@ -1,15 +1,15 @@
package it.cnr.isti.workflow.manager.blocks;
-
import java.util.List;
import java.util.UUID;
-
import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.types.BlockType;
import it.cnr.isti.workflow.manager.bricks.Brick;
+import it.cnr.isti.workflow.manager.executions.executors.BlockExecutor;
import lombok.Builder;
import lombok.Getter;
import lombok.NoArgsConstructor;
+import lombok.NonNull;
import lombok.Singular;
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@@ -18,17 +18,8 @@ public class Block {
final String id = UUID.randomUUID().toString();
- @Builder
- public Block(BlockConfiguration specificConfiguration, String name, T type, @Singular List inputs, @Singular List outputs, @Singular List bricks) {
- this.specificConfiguration = specificConfiguration;
- this.name = name;
- this.type = type;
- this.inputs = inputs;
- this.outputs = outputs;
- }
+ boolean sink;
- T type;
-
String name;
List inputs;
@@ -37,4 +28,22 @@ public class Block {
BlockConfiguration specificConfiguration;
+ BlockExecutor executor;
+
+ BlockType type;
+
+ @Builder
+ public Block(@NonNull BlockConfiguration specificConfiguration, @Singular List inputs,
+ @Singular List outputs, @Singular List bricks, @NonNull BlockExecutor executor, @NonNull T type) {
+ this.specificConfiguration = specificConfiguration;
+ this.name = specificConfiguration.getName();
+ this.inputs = inputs;
+ this.outputs = outputs;
+ this.executor = executor;
+ this.type = type;
+ }
+
+ public void asSink() {
+ this.sink = true;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/BlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/BlockConfiguration.java
index 58ce6e7..09aa08b 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/BlockConfiguration.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/BlockConfiguration.java
@@ -1,6 +1,11 @@
package it.cnr.isti.workflow.manager.blocks.configurations;
+import com.fasterxml.jackson.annotation.JsonIgnore;
+import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
+import com.fasterxml.jackson.annotation.JsonTypeInfo;
+import com.fasterxml.jackson.databind.annotation.JsonTypeIdResolver;
+
import it.cnr.isti.workflow.manager.blocks.types.BlockType;
import jakarta.validation.constraints.NotBlank;
import lombok.Data;
@@ -8,6 +13,13 @@ import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
+@JsonTypeInfo(
+ use = JsonTypeInfo.Id.NAME,
+ include = JsonTypeInfo.As.PROPERTY,
+ property = "type"
+)
+@JsonTypeIdResolver(DynamicBlockConfigurationTypeResolver.class)
+@JsonIgnoreProperties(ignoreUnknown = true)
public abstract class BlockConfiguration {
@NotBlank
@@ -17,8 +29,7 @@ public abstract class BlockConfiguration {
this.name = name;
}
- public abstract Class getType();
-
-
+ @JsonIgnore
+ public abstract Class getBlockType();
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java
new file mode 100644
index 0000000..61d397b
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java
@@ -0,0 +1,59 @@
+package it.cnr.isti.workflow.manager.blocks.configurations;
+
+import com.fasterxml.jackson.annotation.JsonTypeInfo;
+import com.fasterxml.jackson.databind.jsontype.impl.TypeIdResolverBase;
+import com.fasterxml.jackson.databind.JavaType;
+import com.fasterxml.jackson.databind.DatabindContext;
+
+import io.github.classgraph.ClassGraph;
+import io.github.classgraph.ScanResult;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+public class DynamicBlockConfigurationTypeResolver extends TypeIdResolverBase {
+
+ private Map> idToClass = new HashMap<>();
+ private Map, String> classToId = new HashMap<>();
+
+ public DynamicBlockConfigurationTypeResolver() {
+ try (ScanResult scanResult = new ClassGraph().enableClassInfo().scan()) {
+ List> classes = scanResult.getSubclasses(BlockConfiguration.class.getName()).loadClasses();
+ for (Class> clazz : classes) {
+ String id = clazz.getSimpleName(); // puoi scegliere un id più elegante
+ idToClass.put(id, clazz);
+ classToId.put(clazz, id);
+ }
+ }
+ }
+
+ @Override
+ public void init(JavaType baseType) {
+ // Nessuna inizializzazione aggiuntiva necessaria
+ }
+
+ @Override
+ public String idFromValue(Object value) {
+ return classToId.get(value.getClass());
+ }
+
+ @Override
+ public String idFromValueAndType(Object value, Class> suggestedType) {
+ return classToId.get(suggestedType);
+ }
+
+ @Override
+ public JavaType typeFromId(DatabindContext context, String id) {
+ Class> clazz = idToClass.get(id);
+ if (clazz == null) {
+ throw new IllegalArgumentException("Unknown BlockConfiguration type id: " + id);
+ }
+ return context.constructType(clazz);
+ }
+
+ @Override
+ public JsonTypeInfo.Id getMechanism() {
+ return JsonTypeInfo.Id.CUSTOM;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractionBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractionBlockConfiguration.java
index e05b573..79c8efd 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractionBlockConfiguration.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/HumanInteractionBlockConfiguration.java
@@ -1,22 +1,28 @@
package it.cnr.isti.workflow.manager.blocks.configurations;
+import com.fasterxml.jackson.annotation.JsonTypeName;
+
import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType;
+import it.cnr.isti.workflow.manager.bricks.LLMBrick;
import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.NotNull;
import lombok.Data;
import lombok.EqualsAndHashCode;
@Data
@EqualsAndHashCode(callSuper = true)
+@JsonTypeName("HumanInteraction")
public class HumanInteractionBlockConfiguration extends BlockConfiguration {
@NotBlank
private String actionDescription;
+ @NotNull
+ private LLMBrick llmSimulatorBrick;
@Override
- public Class getType() {
+ public Class getBlockType() {
return HumanInteractionBlockType.class;
}
-
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/LLMBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/LLMBlockConfiguration.java
index a004b44..d34bfd6 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/LLMBlockConfiguration.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/LLMBlockConfiguration.java
@@ -30,7 +30,7 @@ public class LLMBlockConfiguration extends BlockConfiguration {
}
@Override
- public Class getType() {
+ public Class getBlockType() {
return LLMBlockType.class;
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/HumanInteractionBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/HumanInteractionBlockFactory.java
new file mode 100644
index 0000000..795bafb
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/HumanInteractionBlockFactory.java
@@ -0,0 +1,35 @@
+package it.cnr.isti.workflow.manager.blocks.factories;
+
+import org.springframework.beans.factory.annotation.Autowired;
+
+import it.cnr.isti.workflow.manager.blocks.Block;
+import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractionBlockConfiguration;
+import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType;
+import it.cnr.isti.workflow.manager.executions.executors.HumanInteractionExecutor;
+import jakarta.validation.Valid;
+
+public class HumanInteractionBlockFactory implements BlockFactory {
+
+ @Autowired
+ HumanInteractionExecutor executor;
+
+ @Autowired
+ HumanInteractionBlockType blockType;
+
+ @Override
+ public Block create(@Valid HumanInteractionBlockConfiguration configuration) {
+ Block block = Block.builder()
+ .input("input").output("output")
+ .specificConfiguration(configuration)
+ .executor(executor)
+ .type(blockType)
+ .build();
+ return block;
+ }
+
+ @Override
+ public Class getBlockType() {
+ return HumanInteractionBlockType.class;
+ }
+
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/LLMBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/LLMBlockFactory.java
index ad08279..ac8391c 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/LLMBlockFactory.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/LLMBlockFactory.java
@@ -9,15 +9,19 @@ import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
+import it.cnr.isti.workflow.manager.executions.executors.LLMExecutor;
import jakarta.validation.Valid;
@Component
public class LLMBlockFactory implements BlockFactory {
+ public static final String OUTPUT_NAME = "response";
@Autowired
LLMBlockType blockType;
+ @Autowired
+ LLMExecutor executor;
@Override
@@ -40,11 +44,14 @@ public class LLMBlockFactory implements BlockFactory create(@Valid LLMBlockConfiguration configuration) {
String prompt =configuration.getPrompt();
Block block = Block.builder()
- .type(blockType).inputs(retrieveInputs(prompt)).output("response")
+ .inputs(retrieveInputs(prompt)).output(OUTPUT_NAME)
.specificConfiguration(configuration)
+ .executor(executor)
+ .type(blockType)
.build();
- return block; }
+ return block;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockType.java
index 57a9c01..b9ee6c0 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockType.java
@@ -1,10 +1,8 @@
package it.cnr.isti.workflow.manager.blocks.types;
-
-
import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration;
-public interface BlockType {
+public interface BlockType {
String getName();
@@ -12,5 +10,8 @@ public interface BlockType {
boolean validate();
+ boolean isUserInteractive();
+
Class extends BlockConfiguration>> getBlockConfigurationClass();
+
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/HumanInteractionBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/HumanInteractionBlockType.java
index ac6816b..6a7a64c 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/HumanInteractionBlockType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/HumanInteractionBlockType.java
@@ -3,6 +3,8 @@ package it.cnr.isti.workflow.manager.blocks.types;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration;
+import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractionBlockConfiguration;
+
@Component(HumanInteractionBlockType.TYPE)
public class HumanInteractionBlockType implements BlockType {
@@ -26,7 +28,11 @@ public class HumanInteractionBlockType implements BlockType {
@Override
public Class extends BlockConfiguration>> getBlockConfigurationClass() {
- return null;
+ return HumanInteractionBlockConfiguration.class;
}
+ @Override
+ public boolean isUserInteractive() {
+ return true;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/LLMBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/LLMBlockType.java
index 8a16fb9..50ad30a 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/LLMBlockType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/LLMBlockType.java
@@ -30,4 +30,8 @@ public class LLMBlockType implements BlockType {
return LLMBlockConfiguration.class;
}
+ @Override
+ public boolean isUserInteractive() {
+ return false;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/SourceBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/SourceBlockType.java
index 455961b..57dc4e7 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/SourceBlockType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/SourceBlockType.java
@@ -29,4 +29,8 @@ public class SourceBlockType implements BlockType {
return null; // Assuming no specific configuration class for SourceBlockType
}
+ @Override
+ public boolean isUserInteractive() {
+ return true;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/bricks/LLMBrick.java b/src/main/java/it/cnr/isti/workflow/manager/bricks/LLMBrick.java
index c2f4c1c..1c84b4f 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/bricks/LLMBrick.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/bricks/LLMBrick.java
@@ -4,13 +4,16 @@ import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
+import lombok.AccessLevel;
+import lombok.NoArgsConstructor;
+@NoArgsConstructor(access = AccessLevel.PRIVATE)
public class LLMBrick extends Brick {
@JsonIgnore
- private final LLMProvider provider;
+ private LLMProvider provider;
- private final String model;
+ private String model;
public String generate(String prompt, Map executionParameters) {
return provider.generate(model, prompt);
diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java
index f04d28b..db39ca1 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java
@@ -53,9 +53,9 @@ public class BlocksController {
@PostMapping
public > Block create(@RequestBody @Valid C blockConfiguration) {
BlockFactory factory = (BlockFactory) blockFactories.stream()
- .filter(f -> f.getBlockType().equals(blockConfiguration.getType()))
+ .filter(f -> f.getBlockType().equals(blockConfiguration.getBlockType()))
.findFirst()
- .orElseThrow(() -> new IllegalArgumentException("Block factory not found for type: " + blockConfiguration.getType()));
+ .orElseThrow(() -> new IllegalArgumentException("Block factory not found for type: " + blockConfiguration.getBlockType()));
Block block = factory.create(blockConfiguration);
return block;
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
index 96d5ac8..2837990 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
@@ -6,49 +6,21 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.logging.Logger;
+
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+import it.cnr.isti.workflow.manager.executions.steps.Step;
+import it.cnr.isti.workflow.manager.executions.steps.StepStatus;
import lombok.Getter;
import lombok.NoArgsConstructor;
@Getter()
@NoArgsConstructor
-public class ExecutionContext {
-
- public enum Status {
- 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;
+public class ExecutionContext implements ExecutionListener {
- private boolean finalState;
- private boolean waitingState;
- Status(boolean initState, boolean finalState, boolean waitingState) {
- this.initState = initState;
- this.finalState = finalState;
- this.waitingState = waitingState;
- }
-
- public boolean isInitState() {
- return initState;
- }
-
- public boolean isFinalState() {
- return finalState;
- }
-
- public boolean isWaitingState() {
- return waitingState ;
- }
-
- public boolean isRunningState() {
- return !finalState && !initState && !waitingState;
- }
- }
+ private static final Logger logger = Logger.getLogger(ExecutionContext.class.getName());
Map inputs = new HashMap<>();
@@ -58,27 +30,26 @@ public class ExecutionContext {
Long startTime = null;
Long endTime = null;
- List stepsUnderExecution = new ArrayList<>();
- List waitingSteps = new ArrayList<>();
Map errors = new HashMap<>();
List warnings = new ArrayList<>();
+ Map> steps;
- Map> nodeResult = new HashMap<>();
+ ExecutionStatus status = ExecutionStatus.CREATED;
- Status status = Status.CREATED;
-
- public void setStatus(Status status) {
- this.status = status;
- if (status == Status.RUNNING) {
- this.startTime = System.currentTimeMillis();
- } else if (status == Status.SUCCESS || status == Status.ERROR) {
- this.endTime = System.currentTimeMillis();
- }
+ public ExecutionContext(Map> steps) {
+ this.steps = steps;
+ this.steps.values().forEach(step -> step.setListener(this));
}
- public void addNodeResult(String nodeId, Map result) {
- this.nodeResult.put(nodeId, result);
+
+ protected void setStatus(ExecutionStatus status) {
+ this.status = status;
+ if (status == ExecutionStatus.RUNNING) {
+ this.startTime = System.currentTimeMillis();
+ } else if (status == ExecutionStatus.SUCCESS || status == ExecutionStatus.ERROR) {
+ this.endTime = System.currentTimeMillis();
+ }
}
public void addResult(String nodeId, String key, Object value) {
@@ -96,5 +67,60 @@ public class ExecutionContext {
public Map getExecutionResult() {
return Collections.unmodifiableMap(this.result);
}
+
+
+ @Override
+ public void completed(String id, Map result) {
+ logger.info("Step " + id + " completed with result: " + result);
+ if (this.steps.get(id).getBlock().isSink()) {
+ logger.info("All inputs for step " + id + " that is a sink have been satisfied.");
+ result.forEach((k, v) -> this.addResult(id, k, v));
+ boolean allCompleted = this.steps.values().stream().filter(s -> s.getBlock().isSink()).allMatch(s -> s.getStatus() == StepStatus.COMPLETED);
+ if (allCompleted) {
+ this.setStatus(ExecutionStatus.SUCCESS);
+ }
+ }
+ }
+
+
+ @Override
+ public void failed(String id, String error) {
+ this.addError(id, error);
+ this.setStatus(ExecutionStatus.ERROR);
+ logger.severe("Step " + id + " failed with error: " + error);
+ }
+
+
+ @Override
+ public void started(String id) {
+ logger.info("Step " + id + " started");
+ }
+
+
+ @Override
+ public void paused(String id) {
+ logger.info("Step " + id + " paused");
+ }
+
+ protected void setInput(String stepId, String inputName, Object value) {
+ Step> step = this.steps.get(stepId);
+ if (step == null)
+ throw new IllegalArgumentException("Step with id " + stepId + " not found");
+ if (step.getStatus() != StepStatus.WAITING_FOR_INPUT && step.getStatus() != StepStatus.READY)
+ throw new IllegalStateException("Step with id " + stepId + " is not in WAITING_FOR_INPUT status (CURRENT STATUS is " + step.getStatus() + ")");
+ Input input = step.getInputs().stream().filter(i -> i.getName().equals(inputName)).findFirst()
+ .orElseThrow(() -> new IllegalArgumentException("Input with name " + inputName + " not found in step with id " + stepId));
+ if (input.isRegistered())
+ throw new IllegalStateException("Input with name " + inputName + " in step with id " + stepId + " is already connected to an output");
+ input.setValue(value);
+
+
+ }
+
+ protected void start(ExecutorService executorService) {
+ this.setStatus(ExecutionStatus.RUNNING);
+ this.steps.values().forEach(step -> step.start(executorService));
+ }
+
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java
new file mode 100644
index 0000000..d926347
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java
@@ -0,0 +1,14 @@
+package it.cnr.isti.workflow.manager.executions;
+
+import java.util.Map;
+
+public interface ExecutionListener {
+
+ void completed(String id, Map result);
+
+ void failed(String id, String error);
+
+ void started(String id);
+
+ void paused(String id);
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
new file mode 100644
index 0000000..50ba540
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
@@ -0,0 +1,95 @@
+package it.cnr.isti.workflow.manager.executions;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.UUID;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import org.bouncycastle.jcajce.provider.asymmetric.dsa.DSASigner.detDSA;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
+
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+import it.cnr.isti.workflow.manager.executions.steps.Output;
+import it.cnr.isti.workflow.manager.executions.steps.Step;
+import it.cnr.isti.workflow.manager.executions.steps.StepStatus;
+import it.cnr.isti.workflow.manager.flows.model.Connection;
+import it.cnr.isti.workflow.manager.flows.model.Flow;
+import lombok.Builder;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+
+@Getter
+@NoArgsConstructor
+public class ExecutionObject {
+
+ String id = UUID.randomUUID().toString();
+
+ ExecutionContext context;
+
+ long creationTime = System.currentTimeMillis();
+
+ String name;
+
+ @JsonIgnore
+ ExecutorService executorService;
+
+ @Builder
+ public ExecutionObject(String executionName, Flow flow) {
+ this.name = executionName;
+
+ List> steps = getStepsFromFlow(flow);
+
+ this.executorService = Executors.newFixedThreadPool(steps.size());
+
+ this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity())));
+
+ }
+
+ List> getStepsFromFlow(Flow flow) {
+ List> steps = new ArrayList<>();
+ flow.getBlocks().forEach(block -> steps.add(new Step<>(block)));
+ for (Connection connection: flow.getConnections()){
+ Step> sourceStep = steps.stream().filter(s -> s.getId().equals(connection.getSourceId())).findFirst().orElseThrow();
+ Step> targetStep = steps.stream().filter(s -> s.getId().equals(connection.getTargetId())).findFirst().orElseThrow();
+ Output output = sourceStep.getOutputs().stream().filter(o -> o.getName().equals(connection.getSourceName())).findFirst().orElseThrow();
+ Input input = targetStep.getInputs().stream().filter(i -> i.getName().equals(connection.getTargetName())).findFirst().orElseThrow();
+ output.registerSource(input);
+ }
+ for (Step> step : steps) {
+ for (Output output : step.getOutputs()) {
+ if (!output.isConnected() && !step.getBlock().isSink())
+ throw new IllegalStateException("Output " + output.getName() + " of step " + step.getId()
+ + " is not connected to any input and the block is not a sink");
+ }
+ }
+ return steps;
+ }
+
+ protected void setInput(String blockId, String inputName, Object input) {
+ //TODO add the interaction step input
+ if (this.context.getStatus().isInitState()){
+ this.context.setInput(blockId, inputName, input);
+ if (this.context.getSteps().values().stream().allMatch(s -> s.getStatus() == StepStatus.READY ||
+ s.getInputs().stream().allMatch(i -> i.isRegistered() || i.isSet()))) {
+ this.context.setStatus(ExecutionStatus.READY);
+ }
+ } else
+ throw new IllegalStateException("Execution with id " + this.getId()
+ + " is not in initialization status (CURRENT STATUS is " + this.getContext().getStatus() + ")");
+ }
+
+ protected void start() {
+ if (this.context.getStatus() == ExecutionStatus.READY){
+ this.context.start(executorService);
+ } else
+ throw new IllegalStateException("Execution with id " + this.getId()
+ + " is not in READY status (CURRENT STATUS is " + this.getContext().getStatus() + ")");
+ }
+
+}
\ No newline at end of file
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionStatus.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionStatus.java
new file mode 100644
index 0000000..a0d8dad
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionStatus.java
@@ -0,0 +1,39 @@
+package it.cnr.isti.workflow.manager.executions;
+
+public enum ExecutionStatus {
+
+ CREATED(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;
+
+ ExecutionStatus(boolean initState, boolean finalState, boolean waitingState) {
+ this.initState = initState;
+ this.finalState = finalState;
+ this.waitingState = waitingState;
+ }
+
+ public boolean isInitState() {
+ return initState;
+ }
+
+ public boolean isFinalState() {
+ return finalState;
+ }
+
+ public boolean isWaitingState() {
+ return waitingState;
+ }
+
+ public boolean isRunningState() {
+ return !finalState && !initState && !waitingState;
+ }
+
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java
new file mode 100644
index 0000000..b68533a
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java
@@ -0,0 +1,64 @@
+package it.cnr.isti.workflow.manager.executions;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+
+import org.springframework.stereotype.Service;
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+import it.cnr.isti.workflow.manager.executions.steps.Output;
+import it.cnr.isti.workflow.manager.executions.steps.Step;
+import it.cnr.isti.workflow.manager.executions.steps.StepStatus;
+import it.cnr.isti.workflow.manager.flows.model.Connection;
+import it.cnr.isti.workflow.manager.flows.model.Flow;
+
+@Service
+public class ExecutionsService {
+
+ private static Map executions = new HashMap<>();
+
+
+ public ExecutionObject createExecution(Flow flow) {
+ ExecutionObject execObject = ExecutionObject.builder()
+ .executionName(flow.getName())
+ .flow(flow)
+ .build();
+ executions.put(execObject.getId(), execObject);
+ return execObject;
+ }
+
+ public ExecutionObject getExecution(String id) {
+ ExecutionObject toReturn = executions.get(id);
+ if (toReturn == null)
+ throw new IllegalArgumentException("Execution with id " + id + " not found");
+ return toReturn;
+ }
+
+ public List getAllExecutions() {
+ return executions.values().stream().toList();
+ }
+
+ public void removeExecution(String id) {
+ if (!executions.containsKey(id))
+ throw new IllegalArgumentException("Execution with id " + id + " not found");
+ if (executions.get(id).getContext().getStatus() == ExecutionStatus.RUNNING)
+ throw new IllegalStateException("Execution with id " + id + " is still running");
+ executions.remove(id);
+ }
+
+ public ExecutionObject prepareInput(String executionId, String blockId, String inputName, Object input) {
+ ExecutionObject eo = getExecution(executionId);
+ eo.setInput(blockId, inputName, input);
+ return eo;
+ }
+
+ public ExecutionObject startExecution(String id) {
+ ExecutionObject eo = getExecution(id);
+ eo.start();
+ return eo;
+ }
+
+
+}
\ No newline at end of file
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/BlockExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/BlockExecutor.java
new file mode 100644
index 0000000..43f45ab
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/BlockExecutor.java
@@ -0,0 +1,25 @@
+package it.cnr.isti.workflow.manager.executions.executors;
+
+import java.util.List;
+import java.util.Map;
+
+import it.cnr.isti.workflow.manager.blocks.Block;
+import it.cnr.isti.workflow.manager.blocks.types.BlockType;
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+
+public interface BlockExecutor {
+
+ Map execute(Block block, List inputs);
+
+ Class getBlockType();
+
+ default boolean isInteractive() {
+ return false;
+ }
+
+ default Map simulate(Block block, List inputs)
+ {
+ throw new UnsupportedOperationException("Simulation not supported for this block type");
+ }
+
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/HumanInteractionExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/HumanInteractionExecutor.java
new file mode 100644
index 0000000..c569802
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/HumanInteractionExecutor.java
@@ -0,0 +1,41 @@
+package it.cnr.isti.workflow.manager.executions.executors;
+
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import it.cnr.isti.workflow.manager.blocks.Block;
+import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractionBlockConfiguration;
+import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType;
+import it.cnr.isti.workflow.manager.bricks.LLMBrick;
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+
+public class HumanInteractionExecutor implements BlockExecutor {
+
+ @Override
+ public Map execute(Block block, List inputs) {
+ throw new UnsupportedOperationException("execute method in human interaction is not supported");
+ }
+
+ @Override
+ public Class getBlockType() {
+ return HumanInteractionBlockType.class;
+ }
+
+
+ @Override
+ public boolean isInteractive() {
+ return true;
+ }
+
+ @Override
+ public Map simulate(Block block, List inputs) {
+ HumanInteractionBlockConfiguration config = (HumanInteractionBlockConfiguration) block.getSpecificConfiguration();
+ String context= inputs.stream().map(input -> input.getName() + "= " + input.getValue()).collect(Collectors.joining(", "));
+ String prompt = "giving the following context as input: " + context + "; perform the following action: " +
+ config.getActionDescription();
+ LLMBrick llmBrick = config.getLlmSimulatorBrick();
+ return Map.of("output", llmBrick.generate(prompt, Map.of()));
+ }
+
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/LLMExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/LLMExecutor.java
new file mode 100644
index 0000000..56589aa
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/LLMExecutor.java
@@ -0,0 +1,38 @@
+package it.cnr.isti.workflow.manager.executions.executors;
+
+import java.util.List;
+import java.util.Map;
+
+import org.springframework.stereotype.Component;
+
+import it.cnr.isti.workflow.manager.blocks.Block;
+import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
+import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory;
+import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
+import it.cnr.isti.workflow.manager.bricks.LLMBrick;
+import it.cnr.isti.workflow.manager.executions.steps.Input;
+
+@Component
+public class LLMExecutor implements BlockExecutor {
+
+ @Override
+ public Map execute(Block block, List inputs) {
+ System.out.println("Executing LLM Block: " + block.getName());
+ LLMBlockConfiguration config = (LLMBlockConfiguration) block.getSpecificConfiguration();
+ String prompt = config.getPrompt();
+
+ for (Input input : inputs) {
+ prompt = prompt.replaceAll("\\$\\{\\{" + input.getName() + "\\}\\}", input.getValue().toString());
+ }
+
+ LLMBrick llmBrick = config.getBrick();
+ String response = llmBrick.generate(prompt, Map.of());
+ return Map.of(LLMBlockFactory.OUTPUT_NAME, response);
+ }
+
+ @Override
+ public Class getBlockType() {
+ return LLMBlockType.class;
+ }
+
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java
index 06e4f8a..1bee0bd 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java
@@ -1,8 +1,10 @@
package it.cnr.isti.workflow.manager.executions.steps;
import lombok.Getter;
+import lombok.ToString;
@Getter
+@ToString
public class Input {
String name;
@@ -10,16 +12,28 @@ public class Input {
private boolean registered = false;
+ private InputListener listener;
+
Input(String name) {
this.name = name;
}
public void setValue(Object value) {
this.value = value;
+ if (listener != null) {
+ listener.onInputSet();
+ }
}
protected void registered() {
this.registered = true;
}
+ public boolean isSet() {
+ return this.value != null;
+ }
+
+ protected void setListener(InputListener listener) {
+ this.listener = listener;
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/InputListener.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/InputListener.java
new file mode 100644
index 0000000..aa1cc06
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/InputListener.java
@@ -0,0 +1,6 @@
+package it.cnr.isti.workflow.manager.executions.steps;
+
+public interface InputListener {
+
+ void onInputSet();
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/LLMStep.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/LLMStep.java
deleted file mode 100644
index 38ba74d..0000000
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/LLMStep.java
+++ /dev/null
@@ -1,15 +0,0 @@
-package it.cnr.isti.workflow.manager.executions.steps;
-
-import it.cnr.isti.workflow.manager.blocks.Block;
-import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
-import lombok.NoArgsConstructor;
-
-@NoArgsConstructor(access = lombok.AccessLevel.PRIVATE)
-public class LLMStep extends Step {
-
- LLMStep(Block block) {
- super(block);
- }
-
-
-}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Output.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Output.java
index c3d4f40..c240803 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Output.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Output.java
@@ -1,5 +1,6 @@
package it.cnr.isti.workflow.manager.executions.steps;
+import java.util.ArrayList;
import java.util.List;
import lombok.Getter;
@@ -9,21 +10,25 @@ public class Output {
@Getter
String name;
- private List consumers = List.of();
+ private List sources = new ArrayList<>();
Output(String name) {
this.name = name;
}
- void register(Input input) {
- this.consumers.add(input);
- input.registered();
+ public void registerSource(Input source) {
+ this.sources.add(source);
+ source.registered();
}
- void setValue(Object value) {
- for (Input consumer : consumers) {
- consumer.setValue(value);
+ public void setValue(Object value) {
+ for (Input source : sources) {
+ source.setValue(value);
}
}
+
+ public boolean isConnected() {
+ return !sources.isEmpty();
+ }
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java
index ab52ebb..314cfcf 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java
@@ -1,32 +1,147 @@
package it.cnr.isti.workflow.manager.executions.steps;
+import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.types.BlockType;
+import it.cnr.isti.workflow.manager.executions.ExecutionListener;
+import lombok.Builder;
import lombok.Getter;
import lombok.NoArgsConstructor;
+import lombok.NonNull;
+import lombok.Setter;
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
-@Getter
-public abstract class Step {
+public class Step implements InputListener {
+ private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(Step.class);
+
+ @Getter
private Block block;
- List inputs = List.of();
- List