diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java index 79484e2..e3781a3 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java @@ -29,7 +29,7 @@ public class DelimitedParserBlockConfiguration extends BlockConfiguration outputs = List.of(); + private List outputs = List.of(); @UiOrder(20) @Structural @@ -37,7 +37,7 @@ public class DelimitedParserBlockConfiguration extends BlockConfiguration outputs) { + public DelimitedParserBlockConfiguration(@NonNull String name, String separator, List outputs) { super(name); this.separator = separator; this.outputs = outputs == null ? List.of() : List.copyOf(outputs); @@ -52,10 +52,16 @@ public class DelimitedParserBlockConfiguration extends BlockConfiguration name != null && !name.isBlank()) .distinct() .count(); long nonBlank = outputs.stream() - .map(SwitchCase::name) + .map(DelimitedParserOutput::name) .filter(name -> name != null && !name.isBlank()) .count(); return distinct == nonBlank; } + + @AssertTrue(message = "a multiple (array) output must be the only declared output") + @JsonIgnore + boolean isMultipleOutputModeValid() { + if (outputs == null || outputs.isEmpty()) { + return true; + } + long multipleCount = outputs.stream().filter(DelimitedParserOutput::multiple).count(); + if (multipleCount == 0) { + return true; + } + return multipleCount == 1 && outputs.size() == 1; + } } \ No newline at end of file diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserOutput.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserOutput.java new file mode 100644 index 0000000..984a98c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserOutput.java @@ -0,0 +1,36 @@ +package it.cnr.isti.workflow.manager.blocks.configurations; + +import com.fasterxml.jackson.annotation.JsonProperty; + +import it.cnr.isti.workflow.manager.ios.IOType; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.Size; + +/** + * One output slot of a {@link DelimitedParserBlockConfiguration}. When + * {@code multiple} is true, this must be the block's only declared output: + * every segment produced by the split goes into it as a list, with no + * fixed-count check. Otherwise it names one fixed-position segment, and the + * block requires exactly as many segments as declared outputs. + */ +public record DelimitedParserOutput( + @NotBlank + @Size(max = 64) + String name, + @JsonProperty(required = false) + IOType type, + @JsonProperty(required = false) + boolean multiple) { + + public DelimitedParserOutput(String name) { + this(name, IOType.TEXT, false); + } + + public DelimitedParserOutput(String name, IOType type) { + this(name, type, false); + } + + public DelimitedParserOutput { + type = type == null ? IOType.TEXT : type; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java index 88cbcad..3c60178 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java @@ -8,7 +8,7 @@ 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.SwitchCase; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserOutput; import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.types.DelimitedParserBlockType; import it.cnr.isti.workflow.manager.ios.IODescriptor; @@ -36,12 +36,13 @@ public class DelimitedParserBlockFactory implements BlockFactoryof() : configuration.getOutputs()) { + for (DelimitedParserOutput outputCase : configuration.getOutputs() == null ? List.of() : configuration.getOutputs()) { if (outputCase == null || outputCase.name() == null || outputCase.name().isBlank()) { continue; } - builder.output(IODescriptor.output(outputCase.name(), outputCase.type(), false, - List.of(new IOCapability(toCapabilityType(outputCase.type()), false)))); + builder.output(IODescriptor.output(outputCase.name(), outputCase.type(), outputCase.multiple(), + List.of(new IOCapability(toCapabilityType(outputCase.type()), false), + new IOCapability(toCapabilityType(outputCase.type()), true)))); } return builder.build(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java index 26d2988..e7d1050 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java @@ -12,24 +12,28 @@ public final class DelimitedResponseParser { } public static List parseParts(String payload, String separator, int expectedParts) { - if (!StringUtils.hasText(payload)) { - throw new IllegalArgumentException("Input payload is empty"); - } - if (!StringUtils.hasText(separator)) { - throw new IllegalArgumentException("Separator is required"); - } if (expectedParts <= 0) { throw new IllegalArgumentException("Expected parts must be greater than zero"); } - - String[] tokens = payload.split(Pattern.quote(separator), -1); - List parts = Arrays.stream(tokens) - .map(String::trim) - .toList(); + List parts = splitAll(payload, separator); if (parts.size() != expectedParts) { throw new IllegalArgumentException( "Expected " + expectedParts + " parts separated by '" + separator + "' but found " + parts.size()); } return parts; } + + /** Splits the payload on the separator with no count check - every segment is kept, trimmed. */ + public static List splitAll(String payload, String separator) { + if (!StringUtils.hasText(payload)) { + throw new IllegalArgumentException("Input payload is empty"); + } + if (!StringUtils.hasText(separator)) { + throw new IllegalArgumentException("Separator is required"); + } + String[] tokens = payload.split(Pattern.quote(separator), -1); + return Arrays.stream(tokens) + .map(String::trim) + .toList(); + } } \ No newline at end of file diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java index 3927ae8..b3d2931 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java @@ -1,20 +1,25 @@ package it.cnr.isti.workflow.manager.executions.executors.blocks; import java.io.File; -import java.util.List; +import java.util.ArrayList; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import org.springframework.stereotype.Component; import org.springframework.util.StringUtils; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; + import it.cnr.isti.workflow.manager.blocks.Block; -import it.cnr.isti.workflow.manager.blocks.configurations.SwitchCase; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserOutput; import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.factories.DelimitedParserBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.DelimitedParserBlockType; import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; +import it.cnr.isti.workflow.manager.executions.NodeExecutionException; import it.cnr.isti.workflow.manager.executions.executors.DelimitedResponseParser; import it.cnr.isti.workflow.manager.executions.steps.Input; import it.cnr.isti.workflow.manager.ios.IOType; @@ -22,6 +27,8 @@ import it.cnr.isti.workflow.manager.ios.IOType; @Component public class DelimitedParserExecutor implements BlockExecutor { + private final ObjectMapper objectMapper = new ObjectMapper(); + @Override public Map execute(Block block, List inputs, Map authorizations, Map executionVariables, @@ -31,14 +38,25 @@ public class DelimitedParserExecutor implements BlockExecutor DelimitedParserBlockFactory.INPUT_NAME.equals(input.getDescriptor().getName())) .findFirst() .map(input -> input.getValue() == null ? null : input.getValue().toString()) - .orElseThrow(() -> new IllegalArgumentException("Missing required input: " + DelimitedParserBlockFactory.INPUT_NAME)); + .orElseThrow(() -> new NodeExecutionException("DELIMITED_PARSER_INPUT_MISSING", + "Missing required input: " + DelimitedParserBlockFactory.INPUT_NAME)); - List outputCases = configuration.getOutputs() == null ? List.of() : configuration.getOutputs(); - List parts = DelimitedResponseParser.parseParts(payload, configuration.getSeparator(), outputCases.size()); + List outputCases = configuration.getOutputs() == null ? List.of() : configuration.getOutputs(); + if (configuration.isArrayMode()) { + DelimitedParserOutput outputCase = outputCases.getFirst(); + List parts = split(payload, configuration.getSeparator()); + List coerced = new ArrayList<>(parts.size()); + for (String part : parts) { + coerced.add(coerce(part, outputCase.type(), outputCase.name())); + } + return Map.of(outputCase.name(), coerced); + } + + List parts = parseFixed(payload, configuration.getSeparator(), outputCases.size()); Map result = new LinkedHashMap<>(); for (int index = 0; index < outputCases.size(); index++) { - SwitchCase outputCase = outputCases.get(index); + DelimitedParserOutput outputCase = outputCases.get(index); result.put(outputCase.name(), coerce(parts.get(index), outputCase.type(), outputCase.name())); } return result; @@ -49,22 +67,61 @@ public class DelimitedParserExecutor implements BlockExecutor split(String payload, String separator) { + try { + return DelimitedResponseParser.splitAll(payload, separator); + } catch (IllegalArgumentException e) { + throw new NodeExecutionException("DELIMITED_PARSER_SPLIT_FAILED", e.getMessage()); + } + } + + private List parseFixed(String payload, String separator, int expectedParts) { + try { + return DelimitedResponseParser.parseParts(payload, separator, expectedParts); + } catch (IllegalArgumentException e) { + throw new NodeExecutionException("DELIMITED_PARSER_SEGMENT_COUNT_MISMATCH", e.getMessage()); + } + } + private Object coerce(String rawValue, IOType type, String outputName) { if (type == IOType.BOOLEAN) { - if ("true".equalsIgnoreCase(rawValue)) { - return Boolean.TRUE; - } - if ("false".equalsIgnoreCase(rawValue)) { - return Boolean.FALSE; - } - throw new IllegalArgumentException("Output " + outputName + " expects a boolean value"); + return coerceBoolean(rawValue, outputName); } if (type == IOType.FILE || type == IOType.CSV) { if (!StringUtils.hasText(rawValue)) { - throw new IllegalArgumentException("Output " + outputName + " expects a non-empty file path"); + throw new NodeExecutionException("DELIMITED_PARSER_INVALID_FILE", + "Output " + outputName + " expects a non-empty file path"); } return new File(rawValue); } + if (type == IOType.JSON) { + return coerceJson(rawValue, outputName); + } return rawValue; } -} \ No newline at end of file + + private Boolean coerceBoolean(String rawValue, String outputName) { + String normalized = rawValue == null ? "" : rawValue.strip().toLowerCase(); + if (normalized.equals("true") || normalized.equals("yes") || normalized.equals("1")) { + return Boolean.TRUE; + } + if (normalized.equals("false") || normalized.equals("no") || normalized.equals("0")) { + return Boolean.FALSE; + } + throw new NodeExecutionException("DELIMITED_PARSER_INVALID_BOOLEAN", + "Output " + outputName + " expects a boolean value (true/false, yes/no, 1/0) but found: " + rawValue); + } + + private JsonNode coerceJson(String rawValue, String outputName) { + if (!StringUtils.hasText(rawValue)) { + throw new NodeExecutionException("DELIMITED_PARSER_INVALID_JSON", + "Output " + outputName + " expects a non-empty JSON segment"); + } + try { + return objectMapper.readTree(rawValue); + } catch (RuntimeException e) { + throw new NodeExecutionException("DELIMITED_PARSER_INVALID_JSON", + "Output " + outputName + " expects a valid JSON segment: " + e.getMessage()); + } + } +} diff --git a/src/main/resources/workflow-editor-init/flows.json b/src/main/resources/workflow-editor-init/flows.json index 3375452..7555ea9 100644 --- a/src/main/resources/workflow-editor-init/flows.json +++ b/src/main/resources/workflow-editor-init/flows.json @@ -6044,5 +6044,182 @@ ], "dependencies": [] } + }, + { + "name": "test delimited parser", + "description": "Delimited-parser example: an LLM produces a pipe-separated assessment (name|assessment|meets-bar), split by DelimitedParserBlock into three typed named outputs (TEXT, TEXT, BOOLEAN).", + "owner": "testuser", + "createdAt": "2026-07-24T12:00:00", + "lastUpdateAt": "2026-07-24T12:00:00", + "published": false, + "finalized": false, + "flow": { + "blocks": [ + { + "id": "ad1f5746-8e35-4325-8de1-fccf80fc5178", + "position": { + "x": 0, + "y": 0 + }, + "name": "collect-candidate-summary", + "inputs": [ + { + "name": "input", + "type": "TEXT", + "multiple": false + } + ], + "outputs": [ + { + "name": "output", + "type": "TEXT", + "multiple": false + } + ], + "specificConfiguration": { + "type": "HumanInteractiveBlockConfiguration", + "name": "collect-candidate-summary", + "actionDescription": "Paste the candidate's name and years of relevant experience." + }, + "typeName": "HumanInteractionBlock" + }, + { + "id": "8c72fa8b-b0e4-4113-8bb0-9e3a0f3c0191", + "position": { + "x": 300, + "y": 0 + }, + "name": "assess-candidate", + "inputs": [ + { + "name": "candidate", + "type": "TEXT", + "multiple": false + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "assess-candidate", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "Given this candidate summary, respond in exactly this format with no extra text: ||\n\nCandidate: ${{candidate}}", + "skills": [] + }, + "typeName": "LLMBlock" + }, + { + "id": "68ab5b69-887e-435c-acb9-11262e477632", + "position": { + "x": 600, + "y": 0 + }, + "name": "parse-assessment", + "inputs": [ + { + "name": "input", + "type": "TEXT", + "multiple": false + } + ], + "outputs": [ + { + "name": "name", + "type": "TEXT", + "multiple": false + }, + { + "name": "assessment", + "type": "TEXT", + "multiple": false + }, + { + "name": "meetsBar", + "type": "BOOLEAN", + "multiple": false + } + ], + "specificConfiguration": { + "type": "DelimitedParserBlockConfiguration", + "name": "parse-assessment", + "separator": "|", + "outputs": [ + { + "name": "name", + "type": "TEXT", + "multiple": false + }, + { + "name": "assessment", + "type": "TEXT", + "multiple": false + }, + { + "name": "meetsBar", + "type": "BOOLEAN", + "multiple": false + } + ] + }, + "typeName": "DelimitedParserBlock" + }, + { + "id": "f1c5273f-0428-459c-9e1a-69554e96e1af", + "position": { + "x": 900, + "y": 0 + }, + "name": "assessment-recorded", + "inputs": [ + { + "name": "input", + "type": "ANY", + "multiple": false + } + ], + "outputs": [], + "specificConfiguration": { + "type": "EndBlockConfiguration", + "name": "assessment-recorded", + "outcomeCode": "TEST_DELIMITED_ASSESSMENT_RECORDED", + "outcomeLabel": "Assessment recorded" + }, + "typeName": "EndBlock" + } + ], + "containers": [], + "connections": [ + { + "id": "fc5834f2-6f75-481f-83a3-40b12f87001b", + "sourceId": "ad1f5746-8e35-4325-8de1-fccf80fc5178", + "sourceName": "output", + "targetId": "8c72fa8b-b0e4-4113-8bb0-9e3a0f3c0191", + "targetName": "candidate" + }, + { + "id": "a7ffcedc-dcb8-42f9-b6a0-0f42c66a4fa7", + "sourceId": "8c72fa8b-b0e4-4113-8bb0-9e3a0f3c0191", + "sourceName": "response", + "targetId": "68ab5b69-887e-435c-acb9-11262e477632", + "targetName": "input" + }, + { + "id": "0be2b2e4-272b-4d3b-ae87-8f96d5590eba", + "sourceId": "68ab5b69-887e-435c-acb9-11262e477632", + "sourceName": "meetsBar", + "targetId": "f1c5273f-0428-459c-9e1a-69554e96e1af", + "targetName": "input" + } + ], + "dependencies": [] + } } ] 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 475e162..24bccd5 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 @@ -33,6 +33,7 @@ import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfi import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.SwitchCase; import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserOutput; 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.LLMBlockType; @@ -880,8 +881,8 @@ public class BlocksControllerTest { .name("Parser") .separator("----") .outputs(List.of( - new SwitchCase("message"), - new SwitchCase("isValid", it.cnr.isti.workflow.manager.ios.IOType.BOOLEAN))) + new DelimitedParserOutput("message"), + new DelimitedParserOutput("isValid", it.cnr.isti.workflow.manager.ios.IOType.BOOLEAN))) .build()); assertNotNull(block); 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 dd70c90..63fc444 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 @@ -35,6 +35,8 @@ import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfigura import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.SwitchCase; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserOutput; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; @@ -64,6 +66,7 @@ import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.MCPAgentBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.MCPAgentChatBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.SwitchBlockFactory; +import it.cnr.isti.workflow.manager.blocks.factories.DelimitedParserBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; @@ -71,6 +74,7 @@ import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; import it.cnr.isti.workflow.manager.blocks.types.MCPAgentBlockType; import it.cnr.isti.workflow.manager.blocks.types.MCPAgentChatBlockType; import it.cnr.isti.workflow.manager.blocks.types.SwitchBlockType; +import it.cnr.isti.workflow.manager.blocks.types.DelimitedParserBlockType; import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.ios.IOType; import it.cnr.isti.workflow.manager.llms.ChatMessage; @@ -166,6 +170,9 @@ public class ExecutionTest { @Autowired SwitchBlockFactory switchBlockFactory; + @Autowired + DelimitedParserBlockFactory delimitedParserBlockFactory; + @Autowired HumanInteractiveBlockFactory humanInteractiveBlockFactory; @@ -1708,6 +1715,62 @@ public class ExecutionTest { && parentExecutionId.equals(iteration.getParentExecutionId()))); } + @Test + public void delimitedParserBlockSplitsIntoFixedNamedOutputsInARealExecution() { + Block parser = delimitedParserBlockFactory.create(DelimitedParserBlockConfiguration.builder() + .name("Parser") + .separator("|") + .outputs(List.of(new DelimitedParserOutput("message"), new DelimitedParserOutput("isValid", IOType.BOOLEAN))) + .build()); + + FlowData flow = FlowData.builder().block(parser).build(); + ExecutionObject execObject = executionsService.createExecution("Delimited parser flow", flow); + + execObject = executionsService.prepareInput(execObject.getId(), parser.getId(), "input", "hello world|true"); + 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()); + } + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + assertEquals("hello world", execObject.getContext().getResult().get(new FieldKey(parser.getId(), "message"))); + assertEquals(true, execObject.getContext().getResult().get(new FieldKey(parser.getId(), "isValid"))); + } + + @Test + public void delimitedParserBlockArrayModeCollectsEverySegmentInARealExecution() { + Block parser = delimitedParserBlockFactory.create(DelimitedParserBlockConfiguration.builder() + .name("Parser") + .separator(",") + .outputs(List.of(new DelimitedParserOutput("items", IOType.TEXT, true))) + .build()); + + FlowData flow = FlowData.builder().block(parser).build(); + ExecutionObject execObject = executionsService.createExecution("Delimited parser array flow", flow); + + execObject = executionsService.prepareInput(execObject.getId(), parser.getId(), "input", "alpha,beta,gamma,delta"); + 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()); + } + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + assertEquals(List.of("alpha", "beta", "gamma", "delta"), + execObject.getContext().getResult().get(new FieldKey(parser.getId(), "items"))); + } + @Test public void containerOutputRejectedByDownstreamInputFailsCleanlyInsteadOfHanging() { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutorTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutorTest.java new file mode 100644 index 0000000..28572ad --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutorTest.java @@ -0,0 +1,150 @@ +package it.cnr.isti.workflow.manager.executions.executors.blocks; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.File; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.Test; + +import tools.jackson.databind.JsonNode; + +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.DelimitedParserOutput; +import it.cnr.isti.workflow.manager.blocks.types.DelimitedParserBlockType; +import it.cnr.isti.workflow.manager.executions.NodeExecutionException; +import it.cnr.isti.workflow.manager.executions.steps.Input; +import it.cnr.isti.workflow.manager.ios.IODescriptor; +import it.cnr.isti.workflow.manager.ios.IOType; + +class DelimitedParserExecutorTest { + + private final DelimitedParserExecutor executor = new DelimitedParserExecutor(); + + @Test + void splitsAFixedNumberOfNamedOutputs() { + Map result = executor.execute( + block("|", List.of(new DelimitedParserOutput("message"), + new DelimitedParserOutput("isValid", IOType.BOOLEAN))), + List.of(input("hello|true")), Map.of(), Map.of(), Map.of(), null); + + assertEquals("hello", result.get("message")); + assertEquals(Boolean.TRUE, result.get("isValid")); + } + + @Test + void rejectsAMismatchedSegmentCount() { + NodeExecutionException exception = assertThrows(NodeExecutionException.class, + () -> executor.execute( + block("|", List.of(new DelimitedParserOutput("a"), new DelimitedParserOutput("b"))), + List.of(input("only-one")), Map.of(), Map.of(), Map.of(), null)); + + assertEquals("DELIMITED_PARSER_SEGMENT_COUNT_MISMATCH", exception.getErrorCode()); + } + + @Test + void rejectsAMissingInput() { + NodeExecutionException exception = assertThrows(NodeExecutionException.class, + () -> executor.execute(block("|", List.of(new DelimitedParserOutput("a"))), + List.of(), Map.of(), Map.of(), Map.of(), null)); + + assertEquals("DELIMITED_PARSER_INPUT_MISSING", exception.getErrorCode()); + } + + @Test + void coercesBooleanFromSeveralAcceptedSpellings() { + Map yes = executor.execute(block("|", List.of(new DelimitedParserOutput("flag", IOType.BOOLEAN))), + List.of(input("yes")), Map.of(), Map.of(), Map.of(), null); + Map zero = executor.execute(block("|", List.of(new DelimitedParserOutput("flag", IOType.BOOLEAN))), + List.of(input("0")), Map.of(), Map.of(), Map.of(), null); + + assertEquals(Boolean.TRUE, yes.get("flag")); + assertEquals(Boolean.FALSE, zero.get("flag")); + } + + @Test + void rejectsAnUnrecognizedBooleanSpelling() { + NodeExecutionException exception = assertThrows(NodeExecutionException.class, + () -> executor.execute(block("|", List.of(new DelimitedParserOutput("flag", IOType.BOOLEAN))), + List.of(input("maybe")), Map.of(), Map.of(), Map.of(), null)); + + assertEquals("DELIMITED_PARSER_INVALID_BOOLEAN", exception.getErrorCode()); + } + + @Test + void coercesFileSegmentsToFile() { + Map result = executor.execute( + block("|", List.of(new DelimitedParserOutput("path", IOType.FILE))), + List.of(input("/tmp/report.txt")), Map.of(), Map.of(), Map.of(), null); + + assertInstanceOf(File.class, result.get("path")); + } + + @Test + void rejectsAnEmptyFileSegment() { + NodeExecutionException exception = assertThrows(NodeExecutionException.class, + () -> executor.execute( + block("|", List.of(new DelimitedParserOutput("label"), new DelimitedParserOutput("path", IOType.FILE))), + List.of(input("x| ")), Map.of(), Map.of(), Map.of(), null)); + + assertEquals("DELIMITED_PARSER_INVALID_FILE", exception.getErrorCode()); + } + + @Test + void coercesAValidJsonSegmentToAJsonNode() { + Map result = executor.execute( + block("~", List.of(new DelimitedParserOutput("payload", IOType.JSON))), + List.of(input("{\"ok\":true}")), Map.of(), Map.of(), Map.of(), null); + + JsonNode node = assertInstanceOf(JsonNode.class, result.get("payload")); + assertTrue(node.path("ok").asBoolean()); + } + + @Test + void rejectsAnInvalidJsonSegment() { + NodeExecutionException exception = assertThrows(NodeExecutionException.class, + () -> executor.execute(block("~", List.of(new DelimitedParserOutput("payload", IOType.JSON))), + List.of(input("{not json")), Map.of(), Map.of(), Map.of(), null)); + + assertEquals("DELIMITED_PARSER_INVALID_JSON", exception.getErrorCode()); + } + + @Test + void arrayModeCollectsEverySegmentRegardlessOfCount() { + Map result = executor.execute( + block(",", List.of(new DelimitedParserOutput("items", IOType.TEXT, true))), + List.of(input("a,b,c,d")), Map.of(), Map.of(), Map.of(), null); + + assertEquals(List.of("a", "b", "c", "d"), result.get("items")); + } + + @Test + void arrayModeCoercesEveryElementToTheDeclaredType() { + Map result = executor.execute( + block(",", List.of(new DelimitedParserOutput("flags", IOType.BOOLEAN, true))), + List.of(input("true,false,yes,no")), Map.of(), Map.of(), Map.of(), null); + + assertEquals(List.of(true, false, true, false), result.get("flags")); + } + + private Input input(String payload) { + return Input.detached(IODescriptor.of("input", IOType.TEXT), payload); + } + + private Block block(String separator, List outputs) { + DelimitedParserBlockConfiguration configuration = DelimitedParserBlockConfiguration.builder() + .name("parser") + .separator(separator) + .outputs(outputs) + .build(); + return Block.builder() + .specificConfiguration(configuration) + .type(new DelimitedParserBlockType()) + .build(); + } +}