added new block types

This commit is contained in:
Lucio Lelii 2025-09-11 17:27:10 +02:00
parent dfd835ba56
commit b485ca2424
34 changed files with 935 additions and 111 deletions

View File

@ -71,6 +71,13 @@
<version>0.11.5</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>io.github.classgraph</groupId>
<artifactId>classgraph</artifactId>
<version>4.8.159</version>
</dependency>
<!-- https://mvnrepository.com/artifact/com.kjetland/mbknor-jackson-jsonschema -->
<dependency>
<groupId>com.kjetland</groupId>

View File

@ -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<T extends BlockType> {
final String id = UUID.randomUUID().toString();
@Builder
public Block(BlockConfiguration<T> specificConfiguration, String name, T type, @Singular List<String> inputs, @Singular List<String> outputs, @Singular List<Brick> bricks) {
this.specificConfiguration = specificConfiguration;
this.name = name;
this.type = type;
this.inputs = inputs;
this.outputs = outputs;
}
boolean sink;
T type;
String name;
List<String> inputs;
@ -37,4 +28,22 @@ public class Block<T extends BlockType> {
BlockConfiguration<T> specificConfiguration;
BlockExecutor<T> executor;
BlockType type;
@Builder
public Block(@NonNull BlockConfiguration<T> specificConfiguration, @Singular List<String> inputs,
@Singular List<String> outputs, @Singular List<Brick> bricks, @NonNull BlockExecutor<T> 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;
}
}

View File

@ -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<T extends BlockType> {
@NotBlank
@ -17,8 +29,7 @@ public abstract class BlockConfiguration<T extends BlockType> {
this.name = name;
}
public abstract Class<T> getType();
@JsonIgnore
public abstract Class<T> getBlockType();
}

View File

@ -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<String, Class<?>> idToClass = new HashMap<>();
private Map<Class<?>, String> classToId = new HashMap<>();
public DynamicBlockConfigurationTypeResolver() {
try (ScanResult scanResult = new ClassGraph().enableClassInfo().scan()) {
List<Class<?>> 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;
}
}

View File

@ -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<HumanInteractionBlockType> {
@NotBlank
private String actionDescription;
@NotNull
private LLMBrick llmSimulatorBrick;
@Override
public Class<HumanInteractionBlockType> getType() {
public Class<HumanInteractionBlockType> getBlockType() {
return HumanInteractionBlockType.class;
}
}

View File

@ -30,7 +30,7 @@ public class LLMBlockConfiguration extends BlockConfiguration<LLMBlockType> {
}
@Override
public Class<LLMBlockType> getType() {
public Class<LLMBlockType> getBlockType() {
return LLMBlockType.class;
}

View File

@ -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<HumanInteractionBlockType, HumanInteractionBlockConfiguration> {
@Autowired
HumanInteractionExecutor executor;
@Autowired
HumanInteractionBlockType blockType;
@Override
public Block<HumanInteractionBlockType> create(@Valid HumanInteractionBlockConfiguration configuration) {
Block<HumanInteractionBlockType> block = Block.<HumanInteractionBlockType>builder()
.input("input").output("output")
.specificConfiguration(configuration)
.executor(executor)
.type(blockType)
.build();
return block;
}
@Override
public Class<HumanInteractionBlockType> getBlockType() {
return HumanInteractionBlockType.class;
}
}

View File

@ -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<LLMBlockType, LLMBlockConfiguration> {
public static final String OUTPUT_NAME = "response";
@Autowired
LLMBlockType blockType;
@Autowired
LLMExecutor executor;
@Override
@ -40,11 +44,14 @@ public class LLMBlockFactory implements BlockFactory<LLMBlockType, LLMBlockConfi
public Block<LLMBlockType> create(@Valid LLMBlockConfiguration configuration) {
String prompt =configuration.getPrompt();
Block<LLMBlockType> block = Block.<LLMBlockType>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;
}
}

View File

@ -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();
}

View File

@ -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;
}
}

View File

@ -30,4 +30,8 @@ public class LLMBlockType implements BlockType {
return LLMBlockConfiguration.class;
}
@Override
public boolean isUserInteractive() {
return false;
}
}

View File

@ -29,4 +29,8 @@ public class SourceBlockType implements BlockType {
return null; // Assuming no specific configuration class for SourceBlockType
}
@Override
public boolean isUserInteractive() {
return true;
}
}

View File

@ -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<String, Object> executionParameters) {
return provider.generate(model, prompt);

View File

@ -53,9 +53,9 @@ public class BlocksController {
@PostMapping
public <T extends BlockType, C extends BlockConfiguration<T>> Block<T> create(@RequestBody @Valid C blockConfiguration) {
BlockFactory<T, C> factory = (BlockFactory<T, C>) 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<T> block = factory.create(blockConfiguration);
return block;
}

View File

@ -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<String, Object> inputs = new HashMap<>();
@ -58,27 +30,26 @@ public class ExecutionContext {
Long startTime = null;
Long endTime = null;
List<String> stepsUnderExecution = new ArrayList<>();
List<String> waitingSteps = new ArrayList<>();
Map<String, String> errors = new HashMap<>();
List<String> warnings = new ArrayList<>();
Map<String, Step<?>> steps;
Map<String, Map<String, Object>> 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<String, Step<?>> steps) {
this.steps = steps;
this.steps.values().forEach(step -> step.setListener(this));
}
public void addNodeResult(String nodeId, Map<String, Object> 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<FieldKey, Object> getExecutionResult() {
return Collections.unmodifiableMap(this.result);
}
@Override
public void completed(String id, Map<String, Object> 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));
}
}

View File

@ -0,0 +1,14 @@
package it.cnr.isti.workflow.manager.executions;
import java.util.Map;
public interface ExecutionListener {
void completed(String id, Map<String, Object> result);
void failed(String id, String error);
void started(String id);
void paused(String id);
}

View File

@ -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<Step<?>> steps = getStepsFromFlow(flow);
this.executorService = Executors.newFixedThreadPool(steps.size());
this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity())));
}
List<Step<?>> getStepsFromFlow(Flow flow) {
List<Step<?>> 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() + ")");
}
}

View File

@ -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;
}
}

View File

@ -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<String, ExecutionObject> 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<ExecutionObject> 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;
}
}

View File

@ -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<B extends BlockType> {
Map<String, Object> execute(Block<B> block, List<Input> inputs);
Class<B> getBlockType();
default boolean isInteractive() {
return false;
}
default Map<String, Object> simulate(Block<B> block, List<Input> inputs)
{
throw new UnsupportedOperationException("Simulation not supported for this block type");
}
}

View File

@ -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<HumanInteractionBlockType> {
@Override
public Map<String, Object> execute(Block<HumanInteractionBlockType> block, List<Input> inputs) {
throw new UnsupportedOperationException("execute method in human interaction is not supported");
}
@Override
public Class<HumanInteractionBlockType> getBlockType() {
return HumanInteractionBlockType.class;
}
@Override
public boolean isInteractive() {
return true;
}
@Override
public Map<String, Object> simulate(Block<HumanInteractionBlockType> block, List<Input> 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()));
}
}

View File

@ -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<LLMBlockType> {
@Override
public Map<String, Object> execute(Block<LLMBlockType> block, List<Input> 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<LLMBlockType> getBlockType() {
return LLMBlockType.class;
}
}

View File

@ -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;
}
}

View File

@ -0,0 +1,6 @@
package it.cnr.isti.workflow.manager.executions.steps;
public interface InputListener {
void onInputSet();
}

View File

@ -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<LLMBlockType> {
LLMStep(Block<LLMBlockType> block) {
super(block);
}
}

View File

@ -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<Input> consumers = List.of();
private List<Input> 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();
}
}

View File

@ -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<B extends BlockType> {
public class Step<B extends BlockType> implements InputListener {
private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(Step.class);
@Getter
private Block<B> block;
List<Input> inputs = List.of();
List<Output> outputs = List.of();
@Getter
private String id;
protected Step(Block<B> block) {
@Getter
private final List<Input> inputs = new ArrayList<>();
@Getter
private final List<Output> outputs = new ArrayList<>();
@Getter
private StepStatus status = StepStatus.WAITING_FOR_INPUT;
@Getter
private boolean started = false;
@JsonIgnore
private ExecutorService executorService;
@JsonIgnore
@Setter
private ExecutionListener listener;
@Setter
@Getter
private boolean simulated = false;
@Builder
public Step(@NonNull Block<B> block) {
// Initialize the step with the provided block
this.block = block;
this.id = block.getId();
this.block.getOutputs().forEach(outputName -> {
Output output = new Output(outputName);
this.outputs.add(output);
});
this.block.getInputs().forEach(inputName -> {
Input input = new Input(inputName);
input.setListener(this);
this.inputs.add(input);
});
}
@Override
public void onInputSet() {
if (this.started && this.status != StepStatus.WAITING_FOR_INPUT) {
throw new IllegalStateException("Step already started. Current status: " + this.status);
}
if (inputs.stream().allMatch(input -> input.isSet())) {
this.status = StepStatus.READY;
if (this.started) {
this.executorService.execute(() -> {
run();
});
}
}
}
public void start(ExecutorService executorService) {
this.started = true;
this.executorService = executorService;
if (this.status == StepStatus.READY) {
this.executorService.execute(() -> {
run();
});
}
}
public void run() {
if (!this.started) {
throw new IllegalStateException("Step not started. Call start() before run().");
}
if (this.status != StepStatus.READY) {
throw new IllegalStateException("Step not ready. Current status: " + this.status);
}
logger.info("Executing step " + this.id + " of block " + this.block.getName());
if (!isSimulated() && !this.getBlock().getType().isUserInteractive()) {
this.status = StepStatus.WAITING_FOR_INTERACTION;
listener.paused(this.id);
} else {
this.status = StepStatus.RUNNING;
listener.started(this.id);
try {
Map<String, Object> outputs = this.block.getExecutor().execute(this.block, this.inputs);
for (Output output : this.outputs) {
if (outputs.containsKey(output.getName())) {
output.setValue(outputs.get(output.getName()));
}
}
this.status = StepStatus.COMPLETED;
listener.completed(this.id, outputs);
} catch (Throwable e) {
this.status = StepStatus.FAILED;
listener.failed(this.id, e.getMessage());
}
logger.info("Finished executing step " + this.id + " of block " + this.block.getName() + " with status "
+ this.status);
}
}
public void interact(Map<String, Object> providedOutputs) {
if (this.getBlock().getType().isUserInteractive() == false) {
throw new IllegalStateException("Step of block " + this.getBlock().getName() + " is not interactive");
}
if (this.status != StepStatus.WAITING_FOR_INTERACTION) {
throw new IllegalStateException("Step not waiting for interaction. Current status: " + this.status);
}
logger.info("Resuming step " + this.id + " of block " + this.block.getName());
this.status = StepStatus.RUNNING;
listener.started(id);
for (Output output : this.outputs) {
if (providedOutputs.containsKey(output.getName())) {
output.setValue(providedOutputs.get(output.getName()));
} else {
throw new IllegalArgumentException("Missing output value for: " + output.getName());
}
}
this.status = StepStatus.COMPLETED;
listener.completed(this.id, providedOutputs);
}
}

View File

@ -0,0 +1,11 @@
package it.cnr.isti.workflow.manager.executions.steps;
public enum StepStatus {
WAITING_FOR_INPUT,
READY,
RUNNING,
WAITING_FOR_INTERACTION,
COMPLETED,
FAILED
}

View File

@ -2,17 +2,28 @@ package it.cnr.isti.workflow.manager.flows.model;
import java.util.UUID;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
public class Connection {
String id = UUID.randomUUID().toString();
String sourceBlockId;
String sourceOutputName;
@Builder
public Connection(String sourceId, String sourceName, String targetId, String targetName) {
this.sourceId = sourceId;
this.sourceName = sourceName;
this.targetId = targetId;
this.targetName = targetName;
}
String targetBlockId;
String targetInputName;
String sourceId;
String sourceName;
String targetId;
String targetName;
}

View File

@ -24,7 +24,7 @@ import lombok.NoArgsConstructor;
@Builder
public class FlowEntity {
@GeneratedValue(strategy = GenerationType.IDENTITY)
@GeneratedValue(strategy = GenerationType.UUID)
@Id
private String id;

View File

@ -4,7 +4,6 @@ import java.util.List;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Profile;
import it.cnr.isti.workflow.manager.bricks.Brick;
import it.cnr.isti.workflow.manager.bricks.BrickManager;

View File

@ -0,0 +1,81 @@
package it.cnr.isti.workflow.manager.executions;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import org.junit.jupiter.api.Test;
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 it.cnr.isti.workflow.manager.MyTestConfiguration;
import it.cnr.isti.workflow.manager.executions.steps.Step;
import it.cnr.isti.workflow.manager.flows.FlowTestCreator;
import it.cnr.isti.workflow.manager.flows.model.Flow;
@SpringBootTest
@Import(MyTestConfiguration.class)
@TestPropertySource(locations = "classpath:test.properties")
public class ExecutionTest {
@Autowired
ExecutionsService executionsService;
@Autowired
FlowTestCreator flowTestCreator;
@Test
public void createExecution() {
Flow flow = flowTestCreator.createFlowWithConnection();
ExecutionObject execObject = executionsService.createExecution(flow);
assertNotNull(execObject);
assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus());
}
@Test
public void createExecutionAndSetInput() {
ExecutionObject eo = createExecutionAndSetInputInternally();
assertEquals(ExecutionStatus.READY, eo.getContext().getStatus());
}
@Test
public void createExecutionSetInputAndStart() {
ExecutionObject eo = createExecutionAndSetInputInternally();
assertEquals(ExecutionStatus.READY, eo.getContext().getStatus());
eo = executionsService.startExecution(eo.getId());
assertEquals(ExecutionStatus.RUNNING, eo.getContext().getStatus());
while (eo.getContext().getStatus() == ExecutionStatus.RUNNING) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
assertTrue(eo.getContext().getStatus().isFinalState());
eo.getContext().getResult().forEach((k,v) -> System.out.println("Result: " + k + " -> " + v));
}
private ExecutionObject createExecutionAndSetInputInternally() {
Flow flow = flowTestCreator.createFlowWithConnection();
ExecutionObject execObject = executionsService.createExecution(flow);
assertNotNull(execObject);
assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus());
for (Step<?> s : execObject.getContext().getSteps().values()){
if (s.getInputs().stream().anyMatch(i -> i.getName().equals("name") && !i.isRegistered())){
executionsService.prepareInput(execObject.getId(), s.getId(), "name", "Frank");
} else if (s.getInputs().stream().anyMatch(i -> i.getName().equals("sister") && !i.isRegistered())){
executionsService.prepareInput(execObject.getId(), s.getId(), "sister", "Anna");
}
}
execObject.getContext().getSteps().values().forEach(s -> {
s.getInputs().forEach(i -> System.out.println("Step " + s.getBlock().getName() + " Input " + i.getName() + " = " + i.getValue()));
});
return executionsService.getExecution(execObject.getId());
}
}

View File

@ -1,10 +1,54 @@
package it.cnr.isti.workflow.manager.flows;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
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.MyTestConfiguration;
import it.cnr.isti.workflow.manager.bricks.BrickManager;
import it.cnr.isti.workflow.manager.bricks.LLMBrick;
import it.cnr.isti.workflow.manager.flows.model.Flow;
@SpringBootTest
@Import(MyTestConfiguration.class)
@TestPropertySource(locations = "classpath:test.properties")
public class FlowTest {
@Autowired
FlowService flowService;
@Autowired
FlowTestCreator flowTestCreator;
@Qualifier("llmBrickTestManager")
@Autowired
BrickManager<LLMBrick> llmBrickManager;
@Test
public void createEmptyFlow() {
Flow flow = Flow.builder().name("Test Flow").description("This is a test flow").build();
flowService.createFlow("testUser", flow );
}
@Test
public void createFlow() {
Flow flow = flowTestCreator.createFlowWithConnection();
flowService.createFlow("testUser", flow );
}
@Test
public void serialize() throws Exception {
Flow flow = flowTestCreator.createFlowWithConnection();
ObjectMapper mapper = new ObjectMapper();
String json = mapper.writerWithDefaultPrettyPrinter().writeValueAsString(flow);
System.out.println(json);
}
}

View File

@ -0,0 +1,59 @@
package it.cnr.isti.workflow.manager.flows;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Import;
import org.springframework.stereotype.Service;
import com.fasterxml.jackson.databind.ObjectMapper;
import it.cnr.isti.workflow.manager.MyTestConfiguration;
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.BrickManager;
import it.cnr.isti.workflow.manager.bricks.LLMBrick;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Flow;
@Service
@Import(MyTestConfiguration.class)
public class FlowTestCreator {
@Qualifier("llmBrickTestManager")
@Autowired
BrickManager<LLMBrick> llmBrickManager;
@Autowired
LLMBlockFactory factory;
public Flow createFlowWithConnection() {
Block<LLMBlockType> block1 = factory.create(LLMBlockConfiguration.builder()
.prompt("Hello, ${{name}}!")
.name("first")
.brick(llmBrickManager.getBricks().get(0))
.build());
Block<LLMBlockType> block2 = factory.create(LLMBlockConfiguration.builder()
.prompt("Hello, ${{name}}! How are you ${{name}}? Is ${{sister}} fine?")
.name("second")
.brick(llmBrickManager.getBricks().get(0))
.build());
block2.asSink();
Connection connection = Connection.builder()
.sourceId(block1.getId())
.sourceName(block1.getOutputs().get(0))
.targetId(block2.getId())
.targetName(block2.getInputs().get(0))
.build();
Flow flow = Flow.builder().name("Test Flow").description("This is a test flow").block(block1).block(block2).connection(connection).build();
return flow;
}
}