diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/IterationBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/IterationBlockConfiguration.java new file mode 100644 index 0000000..2d1b4fe --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/IterationBlockConfiguration.java @@ -0,0 +1,37 @@ +package it.cnr.isti.workflow.manager.blocks.configurations; + +import com.fasterxml.jackson.annotation.JsonProperty; + +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.NonNull; + +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@Getter +@EqualsAndHashCode(callSuper = true) +public class IterationBlockConfiguration extends BlockConfiguration { + + @JsonProperty(required = false) + private String itemVariableName; + + @Builder + public IterationBlockConfiguration(@NonNull String name, String itemVariableName) { + super(name); + this.itemVariableName = itemVariableName; + } + + @Override + public Class getBlockType() { + return IterationBlockType.class; + } + + public static IterationBlockConfiguration empty() { + IterationBlockConfiguration configuration = new IterationBlockConfiguration(); + configuration.name = IterationBlockType.TYPE; + configuration.itemVariableName = "item"; + return configuration; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/BlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/BlockFactory.java index 9c02d12..19fe7bd 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/BlockFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/BlockFactory.java @@ -85,6 +85,7 @@ public interface BlockFactory IOCapabilityType.TEXT; case FILE, CSV -> IOCapabilityType.FILE; + case ANY -> IOCapabilityType.ANY; }; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/IterationBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/IterationBlockFactory.java new file mode 100644 index 0000000..11fb341 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/IterationBlockFactory.java @@ -0,0 +1,59 @@ +package it.cnr.isti.workflow.manager.blocks.factories; + +import java.util.List; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.IOCapability; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; +import it.cnr.isti.workflow.manager.blocks.configurations.IterationBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; +import it.cnr.isti.workflow.manager.ios.IODescriptor; +import it.cnr.isti.workflow.manager.ios.IOType; + +@Component +public class IterationBlockFactory implements BlockFactory { + + public static final String INPUT_NAME = "items"; + public static final String OUTPUT_NAME = "results"; + + private static final List INPUT_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.ANY, true)); + private static final List OUTPUT_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.ANY, true)); + + @Autowired + private IterationBlockType blockType; + + @Override + public Block create(IterationBlockConfiguration configuration) { + return Block.builder() + .input(IODescriptor.input(INPUT_NAME, IOType.ANY, true, INPUT_CAPABILITIES)) + .output(IODescriptor.output(OUTPUT_NAME, IOType.ANY, true, OUTPUT_CAPABILITIES)) + .specificConfiguration(configuration) + .type(blockType) + .build(); + } + + @Override + public Block createEmpty() { + return create(IterationBlockConfiguration.empty()); + } + + @Override + public Class getBlockType() { + return IterationBlockType.class; + } + + @Override + public List supportedInputCapabilities() { + return INPUT_CAPABILITIES; + } + + @Override + public List supportedOutputCapabilities() { + return OUTPUT_CAPABILITIES; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/IterationBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/IterationBlockType.java new file mode 100644 index 0000000..8beee34 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/IterationBlockType.java @@ -0,0 +1,37 @@ +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.IterationBlockConfiguration; + +@Component(IterationBlockType.TYPE) +public class IterationBlockType implements BlockType { + + public static final String TYPE = "IterationBlock"; + + @Override + public String getName() { + return TYPE; + } + + @Override + public String getDescription() { + return "Iterates over an input collection. The current MVP forwards the collection as-is and reserves nested subflow execution for a future version."; + } + + @Override + public boolean validate() { + return true; + } + + @Override + public Class> getBlockConfigurationClass() { + return IterationBlockConfiguration.class; + } + + @Override + public boolean isUserInteractive() { + return false; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/IterationExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/IterationExecutor.java new file mode 100644 index 0000000..824c0f0 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/IterationExecutor.java @@ -0,0 +1,32 @@ +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.factories.IterationBlockFactory; +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; +import it.cnr.isti.workflow.manager.executions.steps.Input; + +@Component +public class IterationExecutor implements BlockExecutor { + + @Override + public Map execute(Block block, List inputs, Map context) { + if (inputs.isEmpty()) { + throw new IllegalArgumentException("IterationBlock requires the items input"); + } + Object value = inputs.getFirst().getValue(); + if (!(value instanceof List values)) { + throw new IllegalArgumentException("IterationBlock expects a list input"); + } + return Map.of(IterationBlockFactory.OUTPUT_NAME, List.copyOf(values)); + } + + @Override + public Class getBlockType() { + return IterationBlockType.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 755170f..2161095 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 @@ -98,6 +98,9 @@ public class Input { throw new IllegalArgumentException("Input " + descriptor.getName() + " expects a file value"); } } + case ANY -> { + // Accept any single runtime value. + } } } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/ios/IODescriptor.java b/src/main/java/it/cnr/isti/workflow/manager/ios/IODescriptor.java index 2b7f073..0bc9530 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/ios/IODescriptor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/ios/IODescriptor.java @@ -67,6 +67,7 @@ public class IODescriptor { return switch (type) { case FILE, CSV -> IOCapabilityType.FILE; case TEXT -> IOCapabilityType.TEXT; + case ANY -> IOCapabilityType.ANY; }; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/ios/IOType.java b/src/main/java/it/cnr/isti/workflow/manager/ios/IOType.java index 2a02227..3f547ae 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/ios/IOType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/ios/IOType.java @@ -6,7 +6,8 @@ import com.fasterxml.jackson.annotation.JsonValue; public enum IOType { TEXT, FILE, - CSV; + CSV, + ANY; @JsonCreator public static IOType fromString(String key) { diff --git a/src/test/java/it/cnr/isti/workflow/manager/blocks/BlockTest.java b/src/test/java/it/cnr/isti/workflow/manager/blocks/BlockTest.java index 190131a..140c606 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/blocks/BlockTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/blocks/BlockTest.java @@ -11,13 +11,16 @@ import org.springframework.test.context.TestPropertySource; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.IterationBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPBridgeBlockConfiguration; import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; import it.cnr.isti.workflow.manager.blocks.factories.BlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.HTTPServerCallBlockFactory; +import it.cnr.isti.workflow.manager.blocks.factories.IterationBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.MCPBridgeBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; import it.cnr.isti.workflow.manager.blocks.types.MCPBridgeBlockType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; @@ -38,6 +41,9 @@ public class BlockTest { @Autowired HTTPServerCallBlockFactory httpServerCallBlockFactory; + @Autowired + IterationBlockFactory iterationBlockFactory; + @Test void createLLMBlock() { @@ -110,4 +116,20 @@ public class BlockTest { assertNotNull(block.getOutputs()); } + @Test + void createIterationBlock() { + BlockFactory factory = iterationBlockFactory; + + IterationBlockConfiguration config = IterationBlockConfiguration.builder() + .name("Iterate candidates") + .itemVariableName("candidate") + .build(); + + Block block = factory.create(config); + assertNotNull(block); + assertNotNull(block.getSpecificConfiguration()); + assertNotNull(block.getInputs()); + assertNotNull(block.getOutputs()); + } + } diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java index 415eada..971da5c 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java @@ -19,10 +19,12 @@ import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.IterationBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPBridgeBlockConfiguration; import it.cnr.isti.workflow.manager.controllers.BlocksController.BlockConfigurationDescriptor; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; import it.cnr.isti.workflow.manager.blocks.types.MCPBridgeBlockType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; @@ -211,6 +213,44 @@ public class BlocksControllerTest { assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("response"))); } + @Test + public void getIterationExampleForType() { + Block block = blocksController.getExampleForType(IterationBlockType.TYPE); + + assertNotNull(block); + assertEquals(IterationBlockType.TYPE, block.getType().getName()); + assertEquals(IterationBlockType.TYPE, block.getName()); + assertNotNull(block.getSpecificConfiguration()); + assertTrue(block.getInputs().stream().anyMatch(input -> input.getName().equals("items"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("results"))); + assertTrue(block.getInputs().stream() + .filter(input -> input.getName().equals("items")) + .findFirst() + .orElseThrow() + .isMultiple()); + assertTrue(block.getInputs().stream() + .filter(input -> input.getName().equals("items")) + .findFirst() + .orElseThrow() + .getValueKinds() + .stream() + .anyMatch(capability -> capability.type() == IOCapabilityType.ANY && capability.multiple())); + } + + @Test + public void createIterationBlock() { + IterationBlockConfiguration configuration = IterationBlockConfiguration.builder() + .name("Iterate") + .itemVariableName("item") + .build(); + + Block block = blocksController.create(configuration); + assertNotNull(block); + assertEquals(IterationBlockType.TYPE, block.getType().getName()); + assertTrue(block.getInputs().stream().anyMatch(input -> input.getName().equals("items"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("results"))); + } + @Test public void httpServerCallSchemaContainsAuthorizationHints() { BlockConfigurationDescriptor descriptor = blocksController diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java index 1e3d52a..d585eca 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java @@ -25,10 +25,13 @@ import it.cnr.isti.workflow.manager.flows.FlowTestCreator; import it.cnr.isti.workflow.manager.flows.model.Flow; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.IterationBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.factories.HTTPServerCallBlockFactory; +import it.cnr.isti.workflow.manager.blocks.factories.IterationBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; +import it.cnr.isti.workflow.manager.blocks.types.IterationBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.ios.IOType; @@ -79,6 +82,9 @@ public class ExecutionTest { @Autowired HTTPServerCallBlockFactory httpServerCallBlockFactory; + @Autowired + IterationBlockFactory iterationBlockFactory; + LLMDescriptor llmBrick = LLMDescriptor.builder() .provider("testProvider") .model("testModel") @@ -207,6 +213,37 @@ public class ExecutionTest { assertEquals(ExecutionStatus.READY, execObject.getContext().getStatus()); } + @Test + public void iterationBlockPassesThroughInputArray() { + Block block = iterationBlockFactory.create(IterationBlockConfiguration.builder() + .name("Iterate values") + .itemVariableName("value") + .build()); + FlowData flow = FlowData.builder().block(block).build(); + + ExecutionObject execObject = executionsService.createExecution("Iteration", flow); + assertEquals(ExecutionStatus.CREATED, execObject.getContext().getStatus()); + + executionsService.prepareInput(execObject.getId(), block.getId(), "items", List.of("a", "b", "c")); + execObject = executionsService.getExecution(execObject.getId()); + assertEquals(ExecutionStatus.READY, execObject.getContext().getStatus()); + + execObject = executionsService.startExecution(execObject.getId()); + while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { + try { + Thread.sleep(50); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + execObject = executionsService.getExecution(execObject.getId()); + } + + assertTrue(execObject.getContext().getStatus().isFinalState()); + Object result = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); + assertEquals(List.of("a", "b", "c"), result); + } + @Test public void singleTextInputRejectsMultipleValues() { Flow flow = flowTestCreator.createFlowwithLLMUnpromptedWithConnection(llmBrick);