From 34c1a60d33388af15bbe2aa8d6735879818341d4 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 26 May 2026 10:15:39 +0200 Subject: [PATCH 1/4] - added a subflow as guard of the loop container - validation of the subflow by type - added type boolean as output of nodes --- .../manager/blocks/IOCapabilityType.java | 1 + .../configurations/JsonSchemaProducer.java | 29 +- .../blocks/configurations/SwitchCase.java | 15 +- .../blocks/factories/BlockFactory.java | 3 + .../ChatInteractionBlockFactory.java | 1 + .../factories/MCPAgentChatBlockFactory.java | 3 + .../blocks/factories/SwitchBlockFactory.java | 16 +- .../retrievers/FlowsFieldRetriever.java | 19 +- .../containers/ContainerSubFlowValidator.java | 17 +- .../ContainerConfiguration.java | 4 + .../LoopContainerConfiguration.java | 180 ++------- .../containers/types/LoopContainerType.java | 2 +- .../ContainerSubFlowValidationRegistry.java | 121 +++++++ .../ContainerSubFlowValidationType.java | 7 + .../controllers/ContainersController.java | 26 +- .../manager/executions/ExecutionObject.java | 1 + .../executions/ExecutionVariableKind.java | 1 + .../manager/executions/ExecutionsService.java | 26 +- .../executors/blocks/SwitchExecutor.java | 22 +- .../containers/GenericContainerExecutor.java | 1 + .../containers/IteratorContainerExecutor.java | 1 + .../containers/LoopContainerExecutor.java | 341 +++++++----------- .../manager/executions/steps/Input.java | 5 + .../flows/validation/FlowDataValidator.java | 66 +++- .../workflow/manager/ios/IODescriptor.java | 1 + .../cnr/isti/workflow/manager/ios/IOType.java | 1 + .../controllers/ContainersControllerTest.java | 135 +++---- .../manager/executions/ExecutionTest.java | 90 +++-- 28 files changed, 618 insertions(+), 517 deletions(-) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationType.java diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/IOCapabilityType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/IOCapabilityType.java index caa628e..527e0cf 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/IOCapabilityType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/IOCapabilityType.java @@ -2,6 +2,7 @@ package it.cnr.isti.workflow.manager.blocks; public enum IOCapabilityType { TEXT, + BOOLEAN, FILE, ANY } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java index 03687aa..77c478a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java @@ -49,13 +49,17 @@ import it.cnr.isti.workflow.manager.configurations.annotations.UiUniqueItemsBy; import jakarta.validation.constraints.Size; import com.fasterxml.jackson.annotation.JsonProperty; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationRegistry; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationType; @Component public class JsonSchemaProducer { private final SchemaGenerator schemaGenerator; + private final ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry; - public JsonSchemaProducer() { + public JsonSchemaProducer(ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry) { + this.containerSubFlowValidationRegistry = containerSubFlowValidationRegistry; SchemaGeneratorConfigBuilder configBuilder = new SchemaGeneratorConfigBuilder( SchemaVersion.DRAFT_7, OptionPreset.PLAIN_JSON) .with(Option.DEFINITIONS_FOR_ALL_OBJECTS) @@ -397,9 +401,32 @@ public class JsonSchemaProducer { if (!retriever.validationUrl().isBlank()) { propertySchema.put("x-retriever-validation-url", retriever.validationUrl()); } + applySubFlowValidationMetadata(propertySchema, ownerClass, entry.getKey()); } } + private void applySubFlowValidationMetadata(ObjectNode propertySchema, Class ownerClass, String propertyName) { + containerSubFlowValidationRegistry.typeForConfigurationField(ownerClass, propertyName) + .ifPresent(validationType -> { + propertySchema.put("x-subflow-validation-type", validationType.name()); + appendSubFlowValidationType(propertySchema, "x-retriever-url", validationType); + appendSubFlowValidationType(propertySchema, "x-retriever-validation-url", validationType); + }); + } + + private void appendSubFlowValidationType(ObjectNode propertySchema, String propertyName, + ContainerSubFlowValidationType validationType) { + JsonNode value = propertySchema.get(propertyName); + if (value == null || !value.isTextual()) { + return; + } + String url = value.asText(); + if (url.contains("type=") || url.contains("validationType=")) { + return; + } + propertySchema.put(propertyName, url + (url.contains("?") ? "&" : "?") + "type=" + validationType.name()); + } + private Map, Map> collectLongTextMetadata(Class rootClass) { Map, Map> result = new HashMap<>(); Set> visited = new HashSet<>(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/SwitchCase.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/SwitchCase.java index 668ecbf..f989bb2 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/SwitchCase.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/SwitchCase.java @@ -1,10 +1,23 @@ 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; public record SwitchCase( @NotBlank @Size(max = 64) - String name) { + String name, + @JsonProperty(required = false) + IOType type) { + + public SwitchCase(String name) { + this(name, IOType.TEXT); + } + + public SwitchCase { + type = type == null ? IOType.TEXT : type; + } } 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 2baecdd..4075754 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,9 @@ public interface BlockFactory IOCapabilityType.FILE; case TEXT -> IOCapabilityType.TEXT; + case BOOLEAN -> IOCapabilityType.BOOLEAN; case ANY -> IOCapabilityType.ANY; }; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/MCPAgentChatBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/MCPAgentChatBlockFactory.java index 8ba5980..eb637ff 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/MCPAgentChatBlockFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/MCPAgentChatBlockFactory.java @@ -129,6 +129,9 @@ public class MCPAgentChatBlockFactory if (type == IOType.TEXT) { return IOCapabilityType.TEXT; } + if (type == IOType.BOOLEAN) { + return IOCapabilityType.BOOLEAN; + } if (type == IOType.FILE || type == IOType.CSV) { return IOCapabilityType.FILE; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/SwitchBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/SwitchBlockFactory.java index 03718c1..495a9cc 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/SwitchBlockFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/SwitchBlockFactory.java @@ -24,7 +24,9 @@ public class SwitchBlockFactory implements BlockFactory INPUT_CAPABILITIES = List.of( new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.TEXT, true)); - private static final List OUTPUT_CAPABILITIES = List.of(new IOCapability(IOCapabilityType.TEXT, false)); + private static final List OUTPUT_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.TEXT, false), + new IOCapability(IOCapabilityType.BOOLEAN, false)); private static final Pattern PLACEHOLDER_PATTERN = Pattern.compile("\\$\\{\\{(.*?)}}"); private static final Pattern SPEL_VARIABLE_PATTERN = Pattern.compile("#([a-zA-Z_][a-zA-Z0-9_]*)"); @@ -43,7 +45,8 @@ public class SwitchBlockFactory implements BlockFactoryof() : configuration.getCases()) { if (outputCase != null && outputCase.name() != null && !outputCase.name().isBlank()) { - builder.output(IODescriptor.output(outputCase.name(), IOType.TEXT, false, OUTPUT_CAPABILITIES)); + builder.output(IODescriptor.output(outputCase.name(), outputCase.type(), false, + List.of(new IOCapability(toCapabilityType(outputCase.type()), false)))); } } return builder.build(); @@ -94,4 +97,13 @@ public class SwitchBlockFactory implements BlockFactory supportedOutputCapabilities() { return OUTPUT_CAPABILITIES; } + + private IOCapabilityType toCapabilityType(IOType type) { + return switch (type) { + case BOOLEAN -> IOCapabilityType.BOOLEAN; + case FILE, CSV -> IOCapabilityType.FILE; + case TEXT -> IOCapabilityType.TEXT; + case ANY -> IOCapabilityType.ANY; + }; + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowsFieldRetriever.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowsFieldRetriever.java index db54522..bc059f1 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowsFieldRetriever.java +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowsFieldRetriever.java @@ -46,19 +46,20 @@ public class FlowsFieldRetriever implements SecureDynamicFieldRetriever { } String context = blankToNull(params.get("context")); + String validationType = blankToNull(params.getOrDefault("type", params.get("validationType"))); boolean validOnly = params.containsKey("validOnly") ? Boolean.parseBoolean(params.get("validOnly")) - : ContainerSubFlowValidator.CONTAINER_CONTEXT.equals(context); + : ContainerSubFlowValidator.CONTAINER_CONTEXT.equals(context) || validationType != null; boolean includeValidation = Boolean.parseBoolean(params.getOrDefault("includeValidation", "false")); return flowRepository.findFlowsByOwnerOrPublic(user.getUsername()).stream() - .map(flow -> toItem(flow, context, includeValidation)) + .map(flow -> toItem(flow, context, validationType, includeValidation)) .filter(item -> !validOnly || item.valid()) .toList(); } - private RetrieverItem toItem(FlowEntity flow, String context, boolean includeValidation) { - List validationErrors = resolveValidationErrors(flow, context); + private RetrieverItem toItem(FlowEntity flow, String context, String validationType, boolean includeValidation) { + List validationErrors = resolveValidationErrors(flow, context, validationType); boolean valid = validationErrors.isEmpty(); List errorsPayload = includeValidation ? validationErrors : List.of(); @@ -78,7 +79,15 @@ public class FlowsFieldRetriever implements SecureDynamicFieldRetriever { errorsPayload); } - private List resolveValidationErrors(FlowEntity flow, String context) { + private List resolveValidationErrors(FlowEntity flow, String context, String validationType) { + if (validationType != null) { + try { + return containerSubFlowValidator.validate(flow.getFlow(), validationType).errors(); + } catch (IllegalArgumentException exception) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, + "Unknown subflow validation type: " + validationType, exception); + } + } if (context == null) { return List.of(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java b/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java index fe84af2..5703350 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java @@ -6,6 +6,8 @@ import java.util.List; import org.springframework.stereotype.Component; import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationRegistry; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationType; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; @@ -17,12 +19,24 @@ public class ContainerSubFlowValidator { public static final String CONTAINER_CONTEXT = "CONTAINER"; private final FlowExecutionValidator flowExecutionValidator; + private final ContainerSubFlowValidationRegistry validationRegistry; - public ContainerSubFlowValidator(FlowExecutionValidator flowExecutionValidator) { + public ContainerSubFlowValidator( + FlowExecutionValidator flowExecutionValidator, + ContainerSubFlowValidationRegistry validationRegistry) { this.flowExecutionValidator = flowExecutionValidator; + this.validationRegistry = validationRegistry; } public ValidationResult validate(FlowData subFlow) { + return validate(subFlow, ContainerSubFlowValidationType.CONTAINER); + } + + public ValidationResult validate(FlowData subFlow, String validationType) { + return validate(subFlow, validationRegistry.resolveType(validationType)); + } + + public ValidationResult validate(FlowData subFlow, ContainerSubFlowValidationType validationType) { List errors = new ArrayList<>(flowExecutionValidator.collectErrors(subFlow)); List openInputs = ContainerFlowInterfaceResolver.getOpenInputs(subFlow); List openOutputs = ContainerFlowInterfaceResolver.getOpenOutputs(subFlow); @@ -50,6 +64,7 @@ public class ContainerSubFlowValidator { "type", "Interactive nodes inside containers are not supported yet"))); } + errors.addAll(validationRegistry.validate(validationType, subFlow, openInputs, openOutputs)); return new ValidationResult( errors.isEmpty(), diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/ContainerConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/ContainerConfiguration.java index 589abf3..3c8e059 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/ContainerConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/ContainerConfiguration.java @@ -7,6 +7,8 @@ import tools.jackson.databind.annotation.JsonTypeIdResolver; import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever; import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription; +import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder; import it.cnr.isti.workflow.manager.containers.types.ContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; @@ -41,6 +43,8 @@ public abstract class ContainerConfiguration { requiresAuth = true, validationUrl = "/containers/validate-subflow") @JsonProperty(required = true) + @UiLabel("Internal Flow") + @UiDescription("this is the flow executed inside the container") FlowData subFlow; public ContainerConfiguration(String name, FlowData subFlow) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java index a241405..287828a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java @@ -1,19 +1,15 @@ package it.cnr.isti.workflow.manager.containers.configurations; -import com.fasterxml.jackson.annotation.JsonAlias; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; -import it.cnr.isti.workflow.manager.configurations.annotations.LongText; +import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever; import it.cnr.isti.workflow.manager.configurations.annotations.Structural; -import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen; +import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription; +import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder; -import it.cnr.isti.workflow.manager.configurations.annotations.UiOptionsFromNode; -import it.cnr.isti.workflow.manager.configurations.annotations.UiRequiredWhen; -import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; -import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import jakarta.validation.Valid; import jakarta.validation.constraints.AssertTrue; import jakarta.validation.constraints.Min; @@ -28,89 +24,37 @@ import lombok.NonNull; @EqualsAndHashCode(callSuper = true) public class LoopContainerConfiguration extends ContainerConfiguration { + public static final String GUARD_OUTPUT = "guard"; + public static final String FEEDBACK_OUTPUT = "feedback"; + public static final String FEEDBACK_INPUT = "feedback"; + @UiOrder(30) @Structural + @Valid + @FieldRetriever( + name = "Flows", + url = "/secure-retriever/Flows/subFlow/items", + structuredData = true, + requiresAuth = true, + validationUrl = "/containers/validate-subflow") + @JsonProperty(required = true) + @UiLabel("Guard Flow") + @UiDescription("this flow is used as guard of the loop, must have as final output a text and a boolean") + private FlowData guardSubFlow; + + @UiOrder(40) + @Structural @Min(1) @JsonProperty(required = true) private Integer maxIterations; - @UiOrder(40) - @JsonProperty(required = true) - @Structural - private boolean useLlm; - - @UiOrder(50) - @UiEnabledWhen(field = "useLlm", equals = "false") - @Structural - @JsonAlias("condition") - @JsonProperty("guardCondition") - @LongText( - placeholder = "Add a deterministic guard condition, for example: ${{response}} == 'done' or ${{outputs.response}} == 'done' or #iteration >= 3", - tip = "Supports SpEL. Inputs and latest exposed outputs are available through placeholders ${{}}. Use outputs.* to reference exposed subflow outputs explicitly. The variable #iteration contains the current 1-based iteration number.", - acceptVariableAsPlaceholder = true) - private String guardCondition; - - @Valid - @UiOrder(60) - @UiEnabledWhen(field = "useLlm", equals = "true", group = "llm") - @UiRequiredWhen(field = "useLlm", equals = "true") - private LLMDescriptor llmDescriptor; - - @UiOrder(70) - @UiEnabledWhen(field = "useLlm", equals = "true", group = "llm") - @UiRequiredWhen(field = "useLlm", equals = "true") - @Structural - @JsonAlias("prompt") - @JsonProperty("guardPrompt") - @LongText( - placeholder = "Add the guard prompt the LLM should use to decide whether the loop should stop", - tip = "The LLM must answer true or false. Inputs, exposed outputs and iteration are included in the evaluation context. Use outputs.* to reference exposed subflow outputs explicitly.", - acceptVariableAsPlaceholder = true) - private String guardPrompt; - - @UiOrder(80) - @Structural - @UiEnabledWhen(field = "subFlow", present = true) - @UiOptionsFromNode(collection = "inputs", valueField = "name", labelField = "name") - @JsonProperty(required = false) - private String feedbackInput; - - @Valid - @UiOrder(90) - @UiEnabledWhen(field = "feedbackPrompt", present = true, group = "feedback") - @UiRequiredWhen(field = "feedbackPrompt", present = true) - private LLMDescriptor feedbackLlmDescriptor; - - @UiOrder(100) - @UiEnabledWhen(field = "feedbackInput", present = true, group = "feedback") - @UiRequiredWhen(field = "feedbackInput", present = true) - @Structural - @JsonProperty(required = false) - @LongText( - placeholder = "Add the prompt that constructs the next value for the selected input when the guard is false", - tip = "Executed only when the loop continues. Use inputs.*, outputs.* and iteration to build the next input value. The LLM response becomes the next value of the selected input.", - acceptVariableAsPlaceholder = true) - private String feedbackPrompt; - @Builder public LoopContainerConfiguration(@NonNull String name, FlowData subFlow, - String guardCondition, - boolean useLlm, - LLMDescriptor llmDescriptor, - String guardPrompt, - String feedbackInput, - LLMDescriptor feedbackLlmDescriptor, - String feedbackPrompt, + FlowData guardSubFlow, Integer maxIterations) { super(name, subFlow); - this.guardCondition = guardCondition; - this.useLlm = useLlm; - this.llmDescriptor = llmDescriptor; - this.guardPrompt = guardPrompt; - this.feedbackInput = feedbackInput; - this.feedbackLlmDescriptor = feedbackLlmDescriptor; - this.feedbackPrompt = feedbackPrompt; + this.guardSubFlow = guardSubFlow == null ? FlowData.builder().build() : guardSubFlow; this.maxIterations = maxIterations == null ? 10 : maxIterations; } @@ -123,89 +67,13 @@ public class LoopContainerConfiguration extends ContainerConfiguration handle.publicName().equals(feedbackInput) && !handle.handle().io().isMultiple()); - } - @AssertTrue(message = "maxIterations must be greater than zero") @JsonIgnore boolean isMaxIterationsValid() { return maxIterations != null && maxIterations > 0; } - - @AssertTrue(message = "LoopContainer subFlow must expose at least one open output") - @JsonIgnore - boolean hasOpenOutputsWhenConfigured() { - return !hasSubFlowNodes() || !ContainerFlowInterfaceResolver.getExposedOutputs(getSubFlow()).isEmpty(); - } - - private boolean hasSubFlowNodes() { - return getSubFlow() != null && !getSubFlow().getNodes().isEmpty(); - } - - @Deprecated - @JsonIgnore - public String getCondition() { - return guardCondition; - } - - @Deprecated - @JsonIgnore - public String getPrompt() { - return guardPrompt; - } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java index 86195fb..24a0120 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java @@ -17,7 +17,7 @@ public class LoopContainerType implements ContainerType { @Override public String getDescription() { - return "A container node that repeatedly executes a subflow until its guard evaluates to true. The guard can be deterministic (SpEL) or evaluated by an LLM."; + return "A container node that repeatedly executes a subflow and then a guard subflow. The guard subflow must expose guard and feedback outputs; guard=true continues the loop with feedback as the next feedback input."; } @Override diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java new file mode 100644 index 0000000..379befa --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java @@ -0,0 +1,121 @@ +package it.cnr.isti.workflow.manager.containers.validation; + +import java.util.EnumMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.validation.ValidationError; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode; +import it.cnr.isti.workflow.manager.ios.IOType; + +@Component +public class ContainerSubFlowValidationRegistry { + + private static final String SUB_FLOW_FIELD = "subFlow"; + private static final String GUARD_SUB_FLOW_FIELD = "guardSubFlow"; + + private final Map rules = + new EnumMap<>(ContainerSubFlowValidationType.class); + + public ContainerSubFlowValidationRegistry() { + rules.put(ContainerSubFlowValidationType.CONTAINER, (subFlow, openInputs, openOutputs) -> List.of()); + rules.put(ContainerSubFlowValidationType.LOOP_BODY, this::validateLoopBody); + rules.put(ContainerSubFlowValidationType.LOOP_GUARD, this::validateLoopGuard); + } + + public ContainerSubFlowValidationType resolveType(String rawType) { + if (rawType == null || rawType.isBlank()) { + return ContainerSubFlowValidationType.CONTAINER; + } + String normalized = rawType.trim() + .replace('-', '_') + .replace('.', '_') + .toUpperCase(Locale.ROOT); + return ContainerSubFlowValidationType.valueOf(normalized); + } + + public Optional typeForConfigurationField(Class configurationClass, + String fieldName) { + if (configurationClass == null || fieldName == null || fieldName.isBlank()) { + return Optional.empty(); + } + if (LoopContainerConfiguration.class.isAssignableFrom(configurationClass) + && SUB_FLOW_FIELD.equals(fieldName)) { + return Optional.of(ContainerSubFlowValidationType.LOOP_BODY); + } + if (LoopContainerConfiguration.class.isAssignableFrom(configurationClass) + && GUARD_SUB_FLOW_FIELD.equals(fieldName)) { + return Optional.of(ContainerSubFlowValidationType.LOOP_GUARD); + } + if (SUB_FLOW_FIELD.equals(fieldName)) { + return Optional.of(ContainerSubFlowValidationType.CONTAINER); + } + return Optional.empty(); + } + + public List validate(ContainerSubFlowValidationType type, FlowData subFlow, + List openInputs, + List openOutputs) { + Rule rule = rules.get(type == null ? ContainerSubFlowValidationType.CONTAINER : type); + return rule == null ? List.of() : rule.validate(subFlow, openInputs, openOutputs); + } + + private List validateLoopBody(FlowData subFlow, + List openInputs, + List openOutputs) { + List exposedOutputs = + ContainerFlowInterfaceResolver.getExposedOutputs(subFlow); + if (exposedOutputs.isEmpty()) { + return List.of(error("subFlow", + "LoopContainer subFlow must expose at least one open output")); + } + boolean hasFeedbackInput = ContainerFlowInterfaceResolver.getExposedInputs(subFlow).stream() + .anyMatch(handle -> LoopContainerConfiguration.FEEDBACK_INPUT.equals(handle.publicName()) + && !handle.handle().io().isMultiple()); + if (hasFeedbackInput) { + return List.of(); + } + return List.of(error("subFlow", + "LoopContainer subFlow must expose an open non-multiple input named feedback")); + } + + private List validateLoopGuard(FlowData subFlow, + List openInputs, + List openOutputs) { + List exposedOutputs = + ContainerFlowInterfaceResolver.getExposedOutputs(subFlow); + boolean hasGuardOutput = hasExposedOutput(exposedOutputs, LoopContainerConfiguration.GUARD_OUTPUT, IOType.BOOLEAN); + boolean hasFeedbackOutput = hasExposedOutput(exposedOutputs, LoopContainerConfiguration.FEEDBACK_OUTPUT, IOType.TEXT); + if (hasGuardOutput && hasFeedbackOutput) { + return List.of(); + } + return List.of(error("guardSubFlow", + "LoopContainer guardSubFlow must expose an open non-multiple boolean output named guard and an open non-multiple text output named feedback")); + } + + private boolean hasExposedOutput(List openOutputs, String name, + IOType... allowedTypes) { + return openOutputs.stream() + .anyMatch(handle -> name.equals(handle.publicName()) + && java.util.Arrays.stream(allowedTypes).anyMatch(type -> type == handle.handle().io().getType()) + && !handle.handle().io().isMultiple()); + } + + private ValidationError error(String field, String message) { + return new ValidationError(ValidationErrorCode.CONTAINER_SUBFLOW_INVALID, "flow", null, field, message); + } + + @FunctionalInterface + private interface Rule { + List validate(FlowData subFlow, + List openInputs, + List openOutputs); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationType.java b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationType.java new file mode 100644 index 0000000..b60ac5a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationType.java @@ -0,0 +1,7 @@ +package it.cnr.isti.workflow.manager.containers.validation; + +public enum ContainerSubFlowValidationType { + CONTAINER, + LOOP_BODY, + LOOP_GUARD +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java index f124e21..6d803aa 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java @@ -12,6 +12,7 @@ import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.server.ResponseStatusException; @@ -73,6 +74,7 @@ public class ContainersController { public static class ContainerSubFlowValidationRequest { private FlowData subFlow; + private String type; public ContainerSubFlowValidationRequest() { } @@ -88,6 +90,14 @@ public class ContainersController { public void setSubFlow(FlowData subFlow) { this.subFlow = subFlow; } + + public String type() { + return type; + } + + public void setType(String type) { + this.type = type; + } } public record ContainerHandleDescriptor(String nodeId, String nodeName, IODescriptor io) { @@ -151,9 +161,17 @@ public class ContainersController { @PostMapping("/validate-subflow") @Operation(summary = "Validate container subflow", description = "Validates a subflow before it is assigned to a container structural configuration.") - public ValidationResult validateContainerSubFlow(@RequestBody ContainerSubFlowValidationRequest request) { + public ValidationResult validateContainerSubFlow(@RequestBody ContainerSubFlowValidationRequest request, + @RequestParam(required = false) String type) { FlowData subFlow = request == null ? null : request.subFlow(); - ContainerSubFlowValidator.ValidationResult validation = containerSubFlowValidator.validate(subFlow); + String validationType = type != null && !type.isBlank() ? type : request == null ? null : request.type(); + ContainerSubFlowValidator.ValidationResult validation; + try { + validation = containerSubFlowValidator.validate(subFlow, validationType); + } catch (IllegalArgumentException exception) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, + "Unknown subflow validation type: " + validationType, exception); + } return new ValidationResult( validation.valid(), @@ -162,6 +180,10 @@ public class ContainersController { validation.openOutputs().stream().map(this::toHandle).toList()); } + public ValidationResult validateContainerSubFlow(ContainerSubFlowValidationRequest request) { + return validateContainerSubFlow(request, null); + } + private ContainerHandleDescriptor toHandle(ContainerFlowInterfaceResolver.OpenHandle handle) { return new ContainerHandleDescriptor(handle.blockId(), handle.blockName(), handle.io()); } 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 index e8cb358..9f205a6 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java @@ -420,6 +420,7 @@ public class ExecutionObject { private ExecutionVariableKind toExecutionVariableKind(IODescriptor descriptor) { return switch (descriptor.getType()) { case TEXT -> ExecutionVariableKind.TEXT; + case BOOLEAN -> ExecutionVariableKind.BOOLEAN; case FILE, CSV -> ExecutionVariableKind.FILE_PATH; case ANY -> ExecutionVariableKind.ANY; }; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java index 7defca3..693c324 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java @@ -3,6 +3,7 @@ package it.cnr.isti.workflow.manager.executions; public enum ExecutionVariableKind { ANY, TEXT, + BOOLEAN, JSON, FILE_PATH, HTTP_RESOURCE, 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 index f37fea9..23da4bf 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -398,9 +398,11 @@ public class ExecutionsService { collectHttpRequirement(requirements, block); } for (Container container : flow.getContainers() == null ? List.>of() : flow.getContainers()) { - collectContainerRequirement(requirements, container); if (container != null && container.getSpecificConfiguration() instanceof ContainerConfiguration containerConfiguration) { collectRequirements(requirements, containerConfiguration.getSubFlow()); + if (containerConfiguration instanceof LoopContainerConfiguration loopConfiguration) { + collectRequirements(requirements, loopConfiguration.getGuardSubFlow()); + } } } } @@ -469,28 +471,6 @@ public class ExecutionsService { .addStepReference(block.getId(), block.getName()); } - private void collectContainerRequirement(Map requirements, Container container) { - if (container == null || !(container.getSpecificConfiguration() instanceof LoopContainerConfiguration configuration)) { - return; - } - if (!configuration.isUseLlm() || configuration.getLlmDescriptor() == null) { - return; - } - LLMDescriptor descriptor = configuration.getLlmDescriptor(); - if (descriptor.provider() == null || descriptor.provider().isBlank()) { - return; - } - LLMProvider provider = resolveProvider(descriptor.provider()); - if (provider == null || !provider.requiresAuthorization()) { - return; - } - String key = provider.authorizationKey(); - requirements.computeIfAbsent(key, - ignored -> new RequirementAccumulator(key, provider.getName(), provider.authorizationFieldName(), - provider.authorizationDescription())) - .addStepReference(container.getId(), container.getName()); - } - private LLMProvider resolveProvider(String providerName) { LLMProvider provider = llmProviders.get(providerName); if (provider != null) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java index f773a99..c219e35 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java @@ -27,6 +27,8 @@ import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport; import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver; import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; 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; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @@ -62,7 +64,7 @@ public class SwitchExecutor implements BlockExecutor { ? evaluateWithLlm(config, inputValues, authorizations, executionVariables, allowedOutputs) : evaluateWithExpression(config.getCondition(), inputValues, executionVariables, allowedOutputs); String payload = resolvePlaceholders(config.getOutputTemplate(), inputValues, executionVariables); - return Map.of(selectedOutput, payload); + return Map.of(selectedOutput, coercePayload(payload, selectedOutput, block.getOutputs())); } @Override @@ -94,6 +96,24 @@ public class SwitchExecutor implements BlockExecutor { return validateSelectedOutput(result.toString(), allowedOutputs); } + private Object coercePayload(String payload, String selectedOutput, List outputs) { + IOType outputType = outputs == null ? IOType.TEXT : outputs.stream() + .filter(output -> selectedOutput.equals(output.getName())) + .map(IODescriptor::getType) + .findFirst() + .orElse(IOType.TEXT); + if (outputType == IOType.BOOLEAN) { + if ("true".equalsIgnoreCase(payload)) { + return Boolean.TRUE; + } + if ("false".equalsIgnoreCase(payload)) { + return Boolean.FALSE; + } + throw new IllegalArgumentException("Switch output " + selectedOutput + " expects a boolean payload"); + } + return payload; + } + private String evaluateWithLlm(SwitchBlockConfiguration config, Map inputValues, Map authorizations, Map executionVariables, Set allowedOutputs) { LLMDescriptor llmDescriptor = config.getLlmDescriptor(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java index 4203b5c..da78173 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java @@ -174,6 +174,7 @@ public class GenericContainerExecutor implements ContainerExecutor it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.TEXT; + case BOOLEAN -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.BOOLEAN; case FILE, CSV -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.FILE_PATH; case ANY -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.ANY; }; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java index 943422e..eef0772 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -238,6 +238,7 @@ public class IteratorContainerExecutor implements ContainerExecutor it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.TEXT; + case BOOLEAN -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.BOOLEAN; case FILE, CSV -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.FILE_PATH; case ANY -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.ANY; }; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java index 03dc2d8..10bcd63 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java @@ -3,18 +3,10 @@ package it.cnr.isti.workflow.manager.executions.executors.containers; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; -import java.util.Objects; -import java.util.regex.Matcher; -import java.util.regex.Pattern; import java.util.stream.Collectors; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.expression.ExpressionParser; -import org.springframework.expression.spel.standard.SpelExpressionParser; -import org.springframework.expression.spel.support.MapAccessor; -import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.stereotype.Component; -import org.springframework.util.StringUtils; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.Timer; @@ -29,45 +21,26 @@ import it.cnr.isti.workflow.manager.executions.ExecutionEventType; import it.cnr.isti.workflow.manager.executions.ExecutionObject; import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport; import it.cnr.isti.workflow.manager.executions.ExecutionStatus; -import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver; import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; import it.cnr.isti.workflow.manager.executions.ExecutionsService; import it.cnr.isti.workflow.manager.executions.FieldKey; import it.cnr.isti.workflow.manager.executions.executors.BooleanLlmResponseParser; import it.cnr.isti.workflow.manager.executions.steps.Input; -import it.cnr.isti.workflow.manager.llms.LLMDescriptor; -import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @Component public class LoopContainerExecutor implements ContainerExecutor { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(LoopContainerExecutor.class); - private static final Pattern PLACEHOLDER_PATTERN = Pattern.compile("\\$\\{\\{(.*?)}}"); private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; - private static final String LLM_SYSTEM_PROMPT = """ - You are a workflow loop evaluator. - Decide whether the loop should stop using the given prompt, inputs, outputs, execution variables and iteration index. - Reply strictly as JSON in the form {"result":true} or {"result":false}. - Do not add any extra text. - """; - private static final String FEEDBACK_SYSTEM_PROMPT = """ - You are a workflow loop feedback constructor. - Generate the next value for the requested input so the next loop iteration can improve the previous result. - Reply only with the raw replacement value for that input. - Do not add markdown, labels, or explanations. - """; private final ExecutionsService executionsService; - private final Map llmProviders; - private final ExpressionParser expressionParser = new SpelExpressionParser(); @Autowired(required = false) MeterRegistry meterRegistry; - public LoopContainerExecutor(ExecutionsService executionsService, Map llmProviders) { + public LoopContainerExecutor(ExecutionsService executionsService) { this.executionsService = executionsService; - this.llmProviders = llmProviders; } @Override @@ -78,6 +51,9 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName = ContainerFlowInterfaceResolver .getExposedInputs(configuration.getSubFlow()).stream() @@ -97,39 +73,18 @@ public class LoopContainerExecutor implements ContainerExecutor globalInputs = ExecutionRuntimeContextSupport.globalView(executionVariables); - executionsService.setGlobalInputDescriptors(innerExecution.getId(), - innerExecution.getRequiredGlobalInputs().stream() - .collect(LinkedHashMap::new, (map, required) -> map.put(required.getName(), - ExecutionVariableDescriptor.builder() - .name(required.getName()) - .kind(mapKind(required)) - .value(globalInputs.get(required.getName())) - .cleanupPolicy(it.cnr.isti.workflow.manager.executions.ExecutionVariableCleanupPolicy.NONE) - .description("Global flow input") - .build()), Map::putAll)); - - for (Map.Entry entry : currentInputs.entrySet()) { - ContainerFlowInterfaceResolver.ExposedHandle exposedHandle = inputPortsByName.get(entry.getKey()); - if (exposedHandle == null) { - throw new IllegalArgumentException("Unknown LoopContainer input port: " + entry.getKey()); - } - ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle(); - executionsService.prepareInput(innerExecution.getId(), handle.blockId(), handle.io().getName(), entry.getValue()); - } - - innerExecution = startAndWait(innerExecution); - executionVariables.clear(); - executionVariables.putAll(innerExecution.getContext().getExecutionVariables()); - executionVariableDescriptors.clear(); - executionVariableDescriptors.putAll(innerExecution.getContext().getExecutionVariableDescriptors()); + configuration.getSubFlow(), + currentInputs, + inputPortsByName, + "LoopContainer", + authorizations, + executionVariables, + executionVariableDescriptors, + eventLogger, + iteration); latestOutputs.clear(); for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedOutputs) { @@ -143,21 +98,24 @@ public class LoopContainerExecutor implements ContainerExecutor iterationInputs = new LinkedHashMap<>(currentInputs); updateInputsForNextIteration(currentInputs, latestOutputs, inputPortsByName); - boolean shouldStop = shouldStop(configuration, currentInputs, latestOutputs, authorizations, executionVariables, iteration); + GuardDecision decision = evaluateGuardSubFlow(configuration, container.getName(), iterationInputs, latestOutputs, + authorizations, executionVariables, executionVariableDescriptors, eventLogger, iteration); if (eventLogger != null) { eventLogger.info(ExecutionEventType.CONTAINER_CONDITION_EVALUATED, "Evaluated loop guard", - Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "stop", shouldStop)); + Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "continue", decision.shouldContinue())); eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, - shouldStop ? "Loop completed at iteration " + iteration : "Completed loop iteration " + iteration, - Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "stop", shouldStop)); + decision.shouldContinue() + ? "Completed loop iteration " + iteration + : "Loop completed at iteration " + iteration, + Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "continue", + decision.shouldContinue())); } - if (shouldStop) { + if (decision.shouldContinue()) { + currentInputs.put(LoopContainerConfiguration.FEEDBACK_INPUT, decision.feedback()); + } else { return new LinkedHashMap<>(latestOutputs); } - - applyFeedbackInputForNextIteration(configuration, currentInputs, iterationInputs, latestOutputs, - authorizations, executionVariables, iteration, eventLogger); } throw new IllegalStateException( @@ -178,121 +136,32 @@ public class LoopContainerExecutor implements ContainerExecutor currentInputs, - Map latestOutputs, Map authorizations, - Map executionVariables, int iteration) { - Map contextValues = buildGuardTemplateValues(currentInputs, latestOutputs, iteration); - return configuration.isUseLlm() - ? evaluateWithLlm(configuration, contextValues, currentInputs, latestOutputs, authorizations, - executionVariables, iteration) - : evaluateWithExpression(configuration.getGuardCondition(), contextValues, currentInputs, latestOutputs, - executionVariables, iteration); - } - - private boolean evaluateWithExpression(String expression, Map values, - Map currentInputs, Map latestOutputs, - Map executionVariables, int iteration) { - String normalizedExpression = normalizeExpression(expression); - StandardEvaluationContext context = new StandardEvaluationContext(values); - context.addPropertyAccessor(new MapAccessor()); - values.forEach(context::setVariable); - context.setVariable("inputs", currentInputs == null ? Map.of() : currentInputs); - context.setVariable("outputs", latestOutputs == null ? Map.of() : latestOutputs); - context.setVariable("vars", executionVariables == null ? Map.of() : executionVariables); - context.setVariable("global", ExecutionRuntimeContextSupport.globalView(executionVariables)); - context.setVariable("context", ExecutionRuntimeContextSupport.contextView(executionVariables)); - context.setVariable("iteration", iteration); - Boolean result = expressionParser.parseExpression(normalizedExpression).getValue(context, Boolean.class); - if (result == null) { - throw new IllegalArgumentException("Loop expression did not resolve to a boolean value"); - } - return result; - } - - private boolean evaluateWithLlm(LoopContainerConfiguration configuration, Map values, - Map currentInputs, Map latestOutputs, - Map authorizations, Map executionVariables, int iteration) { - LLMDescriptor llmDescriptor = configuration.getLlmDescriptor(); - LLMProvider llmProvider = resolveProvider(llmDescriptor.provider()); - ensureAuthorization(llmProvider, llmDescriptor.provider(), authorizations); - - String prompt = buildLlmPrompt(configuration, values, currentInputs, latestOutputs, executionVariables, iteration); - String response = callLlmProvider(llmProvider, llmDescriptor, prompt, authorizations); - try { - return parseBooleanResponse(response); - } catch (IllegalArgumentException e) { - logger.warn("Loop guard LLM returned an unparseable response: provider={}, model={}, iteration={}, response={}", - llmDescriptor.provider(), llmDescriptor.model(), iteration, BooleanLlmResponseParser.preview(response)); - throw e; - } - } - - private void applyFeedbackInputForNextIteration(LoopContainerConfiguration configuration, Map nextInputs, + private GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName, Map iterationInputs, Map latestOutputs, Map authorizations, - Map executionVariables, int iteration, ExecutionEventLogger eventLogger) { - if (!StringUtils.hasText(configuration.getFeedbackPrompt()) || !StringUtils.hasText(configuration.getFeedbackInput())) { - return; - } - - LLMDescriptor llmDescriptor = configuration.getFeedbackLlmDescriptor(); - if (llmDescriptor == null) { - throw new IllegalArgumentException("feedbackLlmDescriptor is required when feedbackPrompt is configured"); - } - - Map values = buildGuardTemplateValues(iterationInputs, latestOutputs, iteration); - String prompt = buildFeedbackPrompt(configuration, values, iterationInputs, latestOutputs, executionVariables, iteration); - LLMProvider llmProvider = resolveProvider(llmDescriptor.provider()); - ensureAuthorization(llmProvider, llmDescriptor.provider(), authorizations); - String response = callLlmProvider(llmProvider, llmDescriptor, prompt, authorizations); - nextInputs.put(configuration.getFeedbackInput(), response); - - if (eventLogger != null) { - eventLogger.info(ExecutionEventType.LLM_REQUEST, - "Constructed next loop input " + configuration.getFeedbackInput(), - Map.of( - "containerType", LoopContainerType.TYPE, - "iteration", iteration, - "purpose", "feedbackInput", - "inputName", configuration.getFeedbackInput(), - "provider", llmDescriptor.provider(), - "model", llmDescriptor.model())); - } - } - - private String buildLlmPrompt(LoopContainerConfiguration configuration, Map values, - Map currentInputs, Map latestOutputs, - Map executionVariables, int iteration) { - String resolvedPrompt = ExecutionTemplateResolver.resolve( - configuration.getGuardPrompt(), - values, - executionVariables == null ? Map.of() : executionVariables); - StringBuilder builder = new StringBuilder(); - builder.append(LLM_SYSTEM_PROMPT).append("\n"); - builder.append("Guard prompt: ").append(resolvedPrompt).append("\n"); - builder.append("Inputs: ").append(currentInputs == null ? Map.of() : currentInputs).append("\n"); - builder.append("Exposed outputs: ").append(latestOutputs == null ? Map.of() : latestOutputs).append("\n"); - builder.append("Current values: ").append(values).append("\n"); - builder.append("Execution variables: ").append(executionVariables == null ? Map.of() : executionVariables).append("\n"); - builder.append("Iteration: ").append(iteration).append("\n"); - return builder.toString(); - } - - private String buildFeedbackPrompt(LoopContainerConfiguration configuration, Map values, - Map iterationInputs, Map latestOutputs, - Map executionVariables, int iteration) { - String resolvedPrompt = ExecutionTemplateResolver.resolve( - configuration.getFeedbackPrompt(), - values, - executionVariables == null ? Map.of() : executionVariables); - StringBuilder builder = new StringBuilder(); - builder.append(FEEDBACK_SYSTEM_PROMPT).append("\n"); - builder.append("Target input: ").append(configuration.getFeedbackInput()).append("\n"); - builder.append("Feedback prompt: ").append(resolvedPrompt).append("\n"); - builder.append("Current iteration inputs: ").append(iterationInputs == null ? Map.of() : iterationInputs).append("\n"); - builder.append("Current iteration exposed outputs: ").append(latestOutputs == null ? Map.of() : latestOutputs).append("\n"); - builder.append("Execution variables: ").append(executionVariables == null ? Map.of() : executionVariables).append("\n"); - builder.append("Iteration: ").append(iteration).append("\n"); - return builder.toString(); + Map executionVariables, + Map executionVariableDescriptors, ExecutionEventLogger eventLogger, + int iteration) { + Map guardInputsByName = ContainerFlowInterfaceResolver + .getExposedInputs(configuration.getGuardSubFlow()).stream() + .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed)); + List guardOutputs = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getGuardSubFlow()); + Map guardInputValues = buildGuardTemplateValues(iterationInputs, latestOutputs, iteration); + ExecutionObject guardExecution = executeSubFlow( + containerName + " guard iteration " + iteration, + configuration.getGuardSubFlow(), + guardInputValues, + guardInputsByName, + "LoopContainerGuard", + authorizations, + executionVariables, + executionVariableDescriptors, + eventLogger, + iteration); + Map guardResult = collectExposedOutputs(guardExecution, guardOutputs); + Object guard = requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT); + Object feedback = requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT); + return new GuardDecision(parseBooleanResponse(String.valueOf(guard)), String.valueOf(feedback)); } private Map buildGuardTemplateValues(Map currentInputs, @@ -310,45 +179,73 @@ public class LoopContainerExecutor implements ContainerExecutor Objects.equals(candidate.getName(), providerName)) - .findFirst() - .orElseThrow(() -> new IllegalArgumentException("Provider not found: " + providerName)); - } - - private void ensureAuthorization(LLMProvider llmProvider, String providerName, Map authorizations) { - String authKey = llmProvider.authorizationKey(); - if (llmProvider.requiresAuthorization() - && (!authorizations.containsKey(authKey) || !StringUtils.hasText(String.valueOf(authorizations.get(authKey))))) { - throw new IllegalArgumentException("Missing authorization for provider: " + providerName); - } - } - - private String callLlmProvider(LLMProvider llmProvider, LLMDescriptor llmDescriptor, String prompt, - Map authorizations) { - String authKey = llmProvider.authorizationKey(); - return llmProvider.requiresAuthorization() - ? llmProvider.generate(llmDescriptor.model(), prompt, String.valueOf(authorizations.get(authKey))) - : llmProvider.generate(llmDescriptor.model(), prompt); - } - - private String normalizeExpression(String expression) { - Matcher matcher = PLACEHOLDER_PATTERN.matcher(expression); - StringBuffer buffer = new StringBuffer(); - while (matcher.find()) { - matcher.appendReplacement(buffer, Matcher.quoteReplacement("#" + matcher.group(1))); - } - matcher.appendTail(buffer); - return buffer.toString(); - } - private boolean parseBooleanResponse(String response) { - return BooleanLlmResponseParser.parse(response, "loop evaluator", "Loop"); + return BooleanLlmResponseParser.parse(response, "loop guard subflow", "Loop"); + } + + private ExecutionObject executeSubFlow(String executionName, it.cnr.isti.workflow.manager.flows.model.FlowData subFlow, + Map inputValues, + Map inputPortsByName, + String containerType, + Map authorizations, + Map executionVariables, + Map executionVariableDescriptors, + ExecutionEventLogger eventLogger, + int iteration) { + ExecutionObject innerExecution = executionsService.createExecution(executionName, subFlow); + forwardInnerEvents(innerExecution, eventLogger, iteration, containerType); + + propagateAuthorizations(innerExecution, authorizations); + executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors); + Map globalInputs = ExecutionRuntimeContextSupport.globalView(executionVariables); + executionsService.setGlobalInputDescriptors(innerExecution.getId(), + innerExecution.getRequiredGlobalInputs().stream() + .collect(LinkedHashMap::new, (map, required) -> map.put(required.getName(), + ExecutionVariableDescriptor.builder() + .name(required.getName()) + .kind(mapKind(required)) + .value(globalInputs.get(required.getName())) + .cleanupPolicy(it.cnr.isti.workflow.manager.executions.ExecutionVariableCleanupPolicy.NONE) + .description("Global flow input") + .build()), Map::putAll)); + + for (Map.Entry entry : inputPortsByName.entrySet()) { + if (!inputValues.containsKey(entry.getKey())) { + throw new IllegalArgumentException( + containerType + " subflow input is missing: " + entry.getKey()); + } + ContainerFlowInterfaceResolver.OpenHandle handle = entry.getValue().handle(); + executionsService.prepareInput(innerExecution.getId(), handle.blockId(), handle.io().getName(), + inputValues.get(entry.getKey())); + } + + innerExecution = startAndWait(innerExecution); + executionVariables.clear(); + executionVariables.putAll(innerExecution.getContext().getExecutionVariables()); + executionVariableDescriptors.clear(); + executionVariableDescriptors.putAll(innerExecution.getContext().getExecutionVariableDescriptors()); + return innerExecution; + } + + private Map collectExposedOutputs(ExecutionObject execution, + List exposedOutputs) { + Map outputs = new LinkedHashMap<>(); + for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedOutputs) { + ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle(); + Object value = execution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); + if (value != null) { + outputs.put(exposedHandle.publicName(), value); + } + } + return outputs; + } + + private Object requireGuardOutput(Map guardResult, String outputName) { + Object value = guardResult.get(outputName); + if (value == null) { + throw new IllegalArgumentException("LoopContainer guardSubFlow must return output: " + outputName); + } + return value; } private void propagateAuthorizations(ExecutionObject innerExecution, Map authorizations) { @@ -431,6 +328,7 @@ public class LoopContainerExecutor implements ContainerExecutor it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.TEXT; + case BOOLEAN -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.BOOLEAN; case FILE, CSV -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.FILE_PATH; case ANY -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.ANY; }; @@ -462,4 +360,7 @@ public class LoopContainerExecutor implements ContainerExecutor { + if (!(value instanceof Boolean)) { + throw new IllegalArgumentException("Input " + descriptor.getName() + " expects a boolean value"); + } + } case FILE, CSV -> { if (!(value instanceof File)) { throw new IllegalArgumentException("Input " + descriptor.getName() + " expects a file value"); diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java index 90ad297..0f6d4aa 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java @@ -20,7 +20,11 @@ import it.cnr.isti.workflow.manager.blocks.types.ConditionalBlockType; import it.cnr.isti.workflow.manager.blocks.types.SwitchBlockType; 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.LoopContainerConfiguration; import it.cnr.isti.workflow.manager.containers.factories.ContainerFactory; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationRegistry; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationType; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.Dependency; import it.cnr.isti.workflow.manager.flows.model.FlowData; @@ -38,6 +42,9 @@ public class FlowDataValidator implements ConstraintValidator> containerFactories; + @Autowired + ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry; + @Override public boolean isValid(FlowData flowData, ConstraintValidatorContext context) { if (flowData == null) { @@ -139,24 +146,9 @@ public class FlowDataValidator implements ConstraintValidator containerConfiguration = container.getSpecificConfiguration(); - FlowData subFlow = containerConfiguration.getSubFlow(); - if (subFlow != null && !subFlow.getNodes().isEmpty()) { - List nestedNodes = subFlow.getNodes(); - if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container)) { - throw validationError(error(ValidationErrorCode.NESTED_CONTAINERS_NOT_SUPPORTED, "container", container.getId(), "specificConfiguration.subFlow", - "Nested containers are not supported")); - } - if (nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) { - throw validationError(error(ValidationErrorCode.INTERACTIVE_NODES_IN_CONTAINER_NOT_SUPPORTED, "container", container.getId(), "specificConfiguration.subFlow", - "Interactive blocks inside containers are not supported yet")); - } - - try { - validateFlowData(subFlow); - } catch (FlowValidationException e) { - throw validationError(error(ValidationErrorCode.CONTAINER_SUBFLOW_INVALID, "container", container.getId(), "specificConfiguration.subFlow", - "Invalid container subFlow: " + e.getMessage())); - } + validateContainerSubFlow(container, containerConfiguration, "subFlow", containerConfiguration.getSubFlow()); + if (containerConfiguration instanceof LoopContainerConfiguration loopConfiguration) { + validateContainerSubFlow(container, containerConfiguration, "guardSubFlow", loopConfiguration.getGuardSubFlow()); } final Container canonicalContainer; @@ -176,6 +168,44 @@ public class FlowDataValidator implements ConstraintValidator container, ContainerConfiguration containerConfiguration, + String fieldName, FlowData subFlow) { + String fieldPath = "specificConfiguration." + fieldName; + if (subFlow == null || subFlow.getNodes().isEmpty()) { + return; + } + + List nestedNodes = subFlow.getNodes(); + if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container)) { + throw validationError(error(ValidationErrorCode.NESTED_CONTAINERS_NOT_SUPPORTED, "container", container.getId(), fieldPath, + "Nested containers are not supported")); + } + if (nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) { + throw validationError(error(ValidationErrorCode.INTERACTIVE_NODES_IN_CONTAINER_NOT_SUPPORTED, "container", container.getId(), fieldPath, + "Interactive blocks inside containers are not supported yet")); + } + + try { + validateFlowData(subFlow); + } catch (FlowValidationException e) { + throw validationError(error(ValidationErrorCode.CONTAINER_SUBFLOW_INVALID, "container", container.getId(), fieldPath, + "Invalid container " + fieldName + ": " + e.getMessage())); + } + + ContainerSubFlowValidationType validationType = containerSubFlowValidationRegistry + .typeForConfigurationField(containerConfiguration.getClass(), fieldName) + .orElse(ContainerSubFlowValidationType.CONTAINER); + List errors = containerSubFlowValidationRegistry.validate( + validationType, + subFlow, + ContainerFlowInterfaceResolver.getOpenInputs(subFlow), + ContainerFlowInterfaceResolver.getOpenOutputs(subFlow)); + if (!errors.isEmpty()) { + ValidationError first = errors.get(0); + throw validationError(error(first.code(), "container", container.getId(), fieldPath, first.message())); + } + } + private void validateConnection(Connection connection, HashMap nodesById) { if (connection == null) { throw validationError(error(ValidationErrorCode.NULL_CONNECTION, "flow", null, "connections", "Flow contains a null connection")); 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 0bc9530..8c2313e 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 BOOLEAN -> IOCapabilityType.BOOLEAN; 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 3f547ae..1e12612 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 @@ -5,6 +5,7 @@ import com.fasterxml.jackson.annotation.JsonValue; public enum IOType { TEXT, + BOOLEAN, FILE, CSV, ANY; diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java index 2c3e3b4..6c20619 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java @@ -19,8 +19,11 @@ import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.Position; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; +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.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; +import it.cnr.isti.workflow.manager.blocks.types.SwitchBlockType; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; @@ -28,7 +31,9 @@ import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfi import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; +import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValidationType; import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.ios.IOType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; @SpringBootTest @@ -44,6 +49,27 @@ public class ContainersControllerTest { @Autowired private RequestMappingHandlerMapping requestMappingHandlerMapping; + private FlowData guardSubFlow() { + Block guardOutput = blocksController.create(SwitchBlockConfiguration.builder() + .name("Guard") + .cases(List.of(new SwitchCase(LoopContainerConfiguration.GUARD_OUTPUT, IOType.BOOLEAN))) + .condition("'" + LoopContainerConfiguration.GUARD_OUTPUT + "'") + .useLlm(false) + .outputTemplate("false") + .build()); + Block feedbackOutput = blocksController.create(SwitchBlockConfiguration.builder() + .name("Feedback") + .cases(List.of(new SwitchCase(LoopContainerConfiguration.FEEDBACK_OUTPUT))) + .condition("'" + LoopContainerConfiguration.FEEDBACK_OUTPUT + "'") + .useLlm(false) + .outputTemplate("next") + .build()); + return FlowData.builder() + .block(guardOutput) + .block(feedbackOutput) + .build(); + } + @Test public void getTypes() { List types = containersController.getTypes(); @@ -82,7 +108,7 @@ public class ContainersControllerTest { .schema(); assertFalse(loopSchema.path("definitions").has("FlowData")); - assertTrue(loopSchema.path("definitions").has("LLMDescriptor")); + assertFalse(loopSchema.path("definitions").has("LLMDescriptor")); assertEquals("#/sharedDefinitions/IODescriptor", catalog.sharedDefinitions().path("FlowData").path("properties").path("globalInputs").path("items").path("$ref").asText()); } @@ -149,10 +175,11 @@ public class ContainersControllerTest { JsonNode schema = (JsonNode) descriptor.schema(); JsonNode subFlow = schema.path("properties").path("subFlow"); assertEquals("Flows", subFlow.path("x-retriever-name").asText()); - assertEquals("/secure-retriever/Flows/subFlow/items", subFlow.path("x-retriever-url").asText()); + assertEquals("/secure-retriever/Flows/subFlow/items?type=CONTAINER", subFlow.path("x-retriever-url").asText()); assertTrue(subFlow.path("x-retriever-structured-data").asBoolean()); assertTrue(subFlow.path("x-retriever-requires-auth").asBoolean()); - assertEquals("/containers/validate-subflow", subFlow.path("x-retriever-validation-url").asText()); + assertEquals("/containers/validate-subflow?type=CONTAINER", subFlow.path("x-retriever-validation-url").asText()); + assertEquals(ContainerSubFlowValidationType.CONTAINER.name(), subFlow.path("x-subflow-validation-type").asText()); assertTrue(subFlow.path("x-ui-structural").asBoolean()); assertTrue(schema.path("required").isArray()); assertTrue(java.util.stream.StreamSupport.stream(schema.path("required").spliterator(), false) @@ -190,7 +217,7 @@ public class ContainersControllerTest { } @Test - public void loopContainerSchemaContainsConditionalLlmPromptMetadata() { + public void loopContainerSchemaContainsGuardSubFlowMetadata() { ContainersController.ContainerConfigurationDescriptor descriptor = containersController.getTypes().stream() .filter(type -> LoopContainerType.TYPE.equals(type.type())) .findFirst() @@ -200,71 +227,34 @@ public class ContainersControllerTest { assertTrue(descriptor.schema() instanceof JsonNode); JsonNode schema = (JsonNode) descriptor.schema(); - JsonNode useLlm = schema.path("properties").path("useLlm"); - JsonNode llmDescriptor = schema.path("properties").path("llmDescriptor"); - JsonNode guardPrompt = schema.path("properties").path("guardPrompt"); - JsonNode guardCondition = schema.path("properties").path("guardCondition"); - JsonNode feedbackInput = schema.path("properties").path("feedbackInput"); - JsonNode feedbackPrompt = schema.path("properties").path("feedbackPrompt"); - JsonNode feedbackLlmDescriptor = schema.path("properties").path("feedbackLlmDescriptor"); + JsonNode subFlow = schema.path("properties").path("subFlow"); + JsonNode guardSubFlow = schema.path("properties").path("guardSubFlow"); + JsonNode maxIterations = schema.path("properties").path("maxIterations"); assertEquals( - List.of("name", "subFlow", "maxIterations", "useLlm", "guardCondition", "llmDescriptor", - "guardPrompt", "feedbackInput", "feedbackLlmDescriptor", "feedbackPrompt", "containerType", - "type"), + List.of("name", "subFlow", "guardSubFlow", "maxIterations", "containerType", "type"), propertyNames(schema.path("properties"))); assertEquals( - List.of("name", "subFlow", "maxIterations", "useLlm", "guardCondition", "llmDescriptor", - "guardPrompt", "feedbackInput", "feedbackLlmDescriptor", "feedbackPrompt", "containerType", - "type"), + List.of("name", "subFlow", "guardSubFlow", "maxIterations", "containerType", "type"), arrayValues(schema.path("x-ui-property-order"))); - assertEquals(4, useLlm.path("x-ui-order").asInt()); - assertEquals(5, guardCondition.path("x-ui-order").asInt()); - - assertTrue(useLlm.path("x-ui-structural").asBoolean()); - assertEquals("textarea", guardCondition.path("x-ui-widget").asText()); - - assertEquals("useLlm", llmDescriptor.path("x-ui-enabled-when").path("field").asText()); - assertEquals("true", llmDescriptor.path("x-ui-enabled-when").path("equals").asText()); - assertEquals("useLlm", llmDescriptor.path("x-ui-required-when").path("field").asText()); - assertEquals("true", llmDescriptor.path("x-ui-required-when").path("equals").asText()); - assertEquals("llm", llmDescriptor.path("x-ui-group").asText()); - - assertTrue(guardPrompt.path("x-ui-structural").asBoolean()); - assertEquals("textarea", guardPrompt.path("x-ui-widget").asText()); - assertEquals("useLlm", guardPrompt.path("x-ui-enabled-when").path("field").asText()); - assertEquals("true", guardPrompt.path("x-ui-enabled-when").path("equals").asText()); - assertEquals("useLlm", guardPrompt.path("x-ui-required-when").path("field").asText()); - assertEquals("true", guardPrompt.path("x-ui-required-when").path("equals").asText()); - assertEquals("llm", guardPrompt.path("x-ui-group").asText()); - - assertEquals("inputs", feedbackInput.path("x-ui-options-from-node").path("collection").asText()); - assertEquals("name", feedbackInput.path("x-ui-options-from-node").path("valueField").asText()); - assertEquals("name", feedbackInput.path("x-ui-options-from-node").path("labelField").asText()); - - assertEquals("feedbackInput", feedbackPrompt.path("x-ui-enabled-when").path("field").asText()); - assertTrue(feedbackPrompt.path("x-ui-required-when").path("present").asBoolean()); - assertEquals("feedback", feedbackPrompt.path("x-ui-group").asText()); - - assertEquals("feedbackPrompt", feedbackLlmDescriptor.path("x-ui-enabled-when").path("field").asText()); - assertTrue(feedbackLlmDescriptor.path("x-ui-required-when").path("present").asBoolean()); - assertEquals("feedback", feedbackLlmDescriptor.path("x-ui-group").asText()); + assertEquals(3, guardSubFlow.path("x-ui-order").asInt()); + assertEquals(4, maxIterations.path("x-ui-order").asInt()); + assertEquals(ContainerSubFlowValidationType.LOOP_BODY.name(), subFlow.path("x-subflow-validation-type").asText()); + assertEquals("/secure-retriever/Flows/subFlow/items?type=LOOP_BODY", subFlow.path("x-retriever-url").asText()); + assertEquals(ContainerSubFlowValidationType.LOOP_GUARD.name(), guardSubFlow.path("x-subflow-validation-type").asText()); + assertEquals("/secure-retriever/Flows/subFlow/items?type=LOOP_GUARD", guardSubFlow.path("x-retriever-url").asText()); + assertEquals("/containers/validate-subflow?type=LOOP_GUARD", guardSubFlow.path("x-retriever-validation-url").asText()); + assertTrue(guardSubFlow.path("x-ui-structural").asBoolean()); } @Test - public void loopContainerAcceptsLegacyConditionAndPromptAliases() throws Exception { + public void loopContainerAcceptsGuardSubFlowConfiguration() throws Exception { String payload = """ { "type": "LoopContainerConfiguration", "name": "Loop", "subFlow": {}, - "condition": "${{response}} == 'done'", - "useLlm": true, - "llmDescriptor": { - "provider": "testProvider", - "model": "testModel" - }, - "prompt": "Stop when ${{outputs.response}} is done", + "guardSubFlow": {}, "maxIterations": 3 } """; @@ -272,8 +262,8 @@ public class ContainersControllerTest { LoopContainerConfiguration configuration = ObjectMapperHolder.mapper.readValue(payload, LoopContainerConfiguration.class); - assertEquals("${{response}} == 'done'", configuration.getGuardCondition()); - assertEquals("Stop when ${{outputs.response}} is done", configuration.getGuardPrompt()); + assertNotNull(configuration.getGuardSubFlow()); + assertEquals(3, configuration.getMaxIterations()); } @Test @@ -402,6 +392,29 @@ public class ContainersControllerTest { .anyMatch(error -> error.message().contains("Interactive nodes inside containers are not supported yet"))); } + @Test + public void validateLoopGuardSubFlowRequiresGuardAndFeedbackOutputs() { + Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("Analyze") + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt("Analyze ${{candidate}}") + .build()); + + ContainersController.ValidationResult invalidResult = containersController.validateContainerSubFlow( + new ContainersController.ContainerSubFlowValidationRequest( + FlowData.builder().block(internalBlock).build()), + ContainerSubFlowValidationType.LOOP_GUARD.name()); + ContainersController.ValidationResult validResult = containersController.validateContainerSubFlow( + new ContainersController.ContainerSubFlowValidationRequest(guardSubFlow()), + ContainerSubFlowValidationType.LOOP_GUARD.name()); + + assertFalse(invalidResult.valid()); + assertTrue(invalidResult.errors().stream() + .anyMatch(error -> error.message().contains("guardSubFlow must expose an open non-multiple boolean output named guard"))); + assertTrue(validResult.valid()); + assertTrue(validResult.errors().isEmpty()); + } + @Test public void createGenericContainerAllowsEmptySubFlow() { Container container = containersController.create(GenericContainerConfiguration.builder() @@ -510,19 +523,19 @@ public class ContainersControllerTest { .provider("testProvider") .model("testModel") .build()) - .prompt("Analyze ${{candidate}}") + .prompt("Analyze ${{candidate}} with ${{feedback}}") .build()); Container container = containersController.create(LoopContainerConfiguration.builder() .name("Loop") .subFlow(FlowData.builder().block(internalBlock).build()) - .guardCondition("${{response}} == 'done'") - .useLlm(false) + .guardSubFlow(guardSubFlow()) .maxIterations(3) .build()); assertNotNull(container); assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate"))); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("feedback"))); assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response"))); } 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 df2294e..d61c518 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 @@ -116,6 +116,12 @@ public class ExecutionTest { if (prompt.contains("Construct the next request for the function loop")) { return "Implement the same function and add null validation"; } + if (prompt.contains("Loop guard subflow")) { + return prompt.contains("function v1") ? "true" : "false"; + } + if (prompt.contains("Loop feedback subflow")) { + return "Implement the same function and add null validation"; + } if (prompt.contains("Loop guard via exposed outputs")) { return prompt.contains("Hello, Alice!") ? "{\"result\":true}" : "{\"result\":false}"; } @@ -182,6 +188,52 @@ public class ExecutionTest { .model("testModel") .build(); + private FlowData loopGuardSubFlow() { + Block guardEvaluator = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Guard Evaluator") + .llmDescriptor(llmBrick) + .prompt("Loop guard subflow. Previous output: ${{outputs.response}}") + .build()); + Block guardOutput = switchBlockFactory.create(SwitchBlockConfiguration.builder() + .name("Expose Guard") + .cases(List.of(new SwitchCase(LoopContainerConfiguration.GUARD_OUTPUT, IOType.BOOLEAN))) + .condition("'" + LoopContainerConfiguration.GUARD_OUTPUT + "'") + .useLlm(false) + .outputTemplate("${{response}}") + .build()); + Block feedbackBuilder = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Feedback Builder") + .llmDescriptor(llmBrick) + .prompt("Loop feedback subflow.") + .build()); + Block feedbackOutput = switchBlockFactory.create(SwitchBlockConfiguration.builder() + .name("Expose Feedback") + .cases(List.of(new SwitchCase(LoopContainerConfiguration.FEEDBACK_OUTPUT))) + .condition("'" + LoopContainerConfiguration.FEEDBACK_OUTPUT + "'") + .useLlm(false) + .outputTemplate("${{response}}") + .build()); + + return FlowData.builder() + .block(guardEvaluator) + .block(guardOutput) + .block(feedbackBuilder) + .block(feedbackOutput) + .connection(Connection.builder() + .sourceId(guardEvaluator.getId()) + .sourceName(LLMBlockFactory.OUTPUT_NAME) + .targetId(guardOutput.getId()) + .targetName(LLMBlockFactory.OUTPUT_NAME) + .build()) + .connection(Connection.builder() + .sourceId(feedbackBuilder.getId()) + .sourceName(LLMBlockFactory.OUTPUT_NAME) + .targetId(feedbackOutput.getId()) + .targetName(LLMBlockFactory.OUTPUT_NAME) + .build()) + .build(); + } + @Test public void createExecution() { Flow flow = flowTestCreator.createFlowWithConnection(llmBrick); @@ -1277,21 +1329,20 @@ public class ExecutionTest { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() .name("Internal LLM") .llmDescriptor(llmBrick) - .prompt("Hello, ${{name}}!") + .prompt("Hello, ${{feedback}}!") .build()); Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() .name("Loop") .subFlow(FlowData.builder().block(internalBlock).build()) - .guardCondition("${{response}} == 'Hello, Alice!'") - .useLlm(false) + .guardSubFlow(loopGuardSubFlow()) .maxIterations(3) .build()); FlowData flow = FlowData.builder().container(container).build(); ExecutionObject execObject = executionsService.createExecution("Loop flow", flow); - executionsService.prepareInput(execObject.getId(), container.getId(), "name", "Alice"); + executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Alice"); execObject = executionsService.startExecution(execObject.getId()); while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { try { @@ -1319,21 +1370,20 @@ public class ExecutionTest { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() .name("Internal LLM") .llmDescriptor(llmBrick) - .prompt("Hello, ${{name}}!") + .prompt("Hello, ${{feedback}}!") .build()); Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() .name("Loop") .subFlow(FlowData.builder().block(internalBlock).build()) - .guardCondition("${{outputs.response}} == 'Hello, Alice!'") - .useLlm(false) + .guardSubFlow(loopGuardSubFlow()) .maxIterations(3) .build()); FlowData flow = FlowData.builder().container(container).build(); ExecutionObject execObject = executionsService.createExecution("Loop outputs flow", flow); - executionsService.prepareInput(execObject.getId(), container.getId(), "name", "Alice"); + executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Alice"); execObject = executionsService.startExecution(execObject.getId()); while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { execObject = executionsService.getExecution(execObject.getId()); @@ -1349,22 +1399,20 @@ public class ExecutionTest { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() .name("Internal LLM") .llmDescriptor(llmBrick) - .prompt("Hello, ${{name}}!") + .prompt("Hello, ${{feedback}}!") .build()); Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() .name("Loop") .subFlow(FlowData.builder().block(internalBlock).build()) - .useLlm(true) - .llmDescriptor(llmBrick) - .guardPrompt("Loop guard via exposed outputs: stop when ${{outputs.response}} is exactly Hello, Alice!") + .guardSubFlow(loopGuardSubFlow()) .maxIterations(3) .build()); FlowData flow = FlowData.builder().container(container).build(); ExecutionObject execObject = executionsService.createExecution("Loop llm outputs flow", flow); - executionsService.prepareInput(execObject.getId(), container.getId(), "name", "Alice"); + executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Alice"); execObject = executionsService.startExecution(execObject.getId()); while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { execObject = executionsService.getExecution(execObject.getId()); @@ -1380,25 +1428,20 @@ public class ExecutionTest { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() .name("Function Builder") .llmDescriptor(llmBrick) - .prompt("Build function for request: ${{request}}") + .prompt("Build function for request: ${{feedback}}") .build()); Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() .name("Loop") .subFlow(FlowData.builder().block(internalBlock).build()) - .guardCondition("${{outputs.response}} == 'function v2'") - .feedbackInput("request") - .feedbackLlmDescriptor(llmBrick) - .feedbackPrompt( - "Construct the next request for the function loop. Previous request: ${{inputs.request}} Previous output: ${{outputs.response}} Iteration: ${{iteration}}") - .useLlm(false) + .guardSubFlow(loopGuardSubFlow()) .maxIterations(3) .build()); FlowData flow = FlowData.builder().container(container).build(); ExecutionObject execObject = executionsService.createExecution("Loop feedback flow", flow); - executionsService.prepareInput(execObject.getId(), container.getId(), "request", "Implement the same function"); + executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Implement the same function"); execObject = executionsService.startExecution(execObject.getId()); while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { execObject = executionsService.getExecution(execObject.getId()); @@ -1408,10 +1451,7 @@ public class ExecutionTest { Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); assertEquals("function v2", output); assertTrue(execObject.getContext().getEvents().stream() - .anyMatch(event -> event.getType() == ExecutionEventType.LLM_REQUEST - && "LoopContainer".equals(event.getDetails().get("containerType")) - && "feedbackInput".equals(event.getDetails().get("purpose")) - && "request".equals(event.getDetails().get("inputName")))); + .anyMatch(event -> "LoopContainerGuard".equals(event.getDetails().get("containerType")))); } @Test From 83c76ebd4f02acddba3f3c270522f2ce277691bf Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 26 May 2026 23:13:46 +0200 Subject: [PATCH 2/4] Add generic DelimitedParser block with configurable outputs --- .../DelimitedParserBlockConfiguration.java | 88 +++++++++++++++++++ .../DelimitedParserBlockFactory.java | 77 ++++++++++++++++ .../types/DelimitedParserBlockType.java | 37 ++++++++ .../executors/DelimitedResponseParser.java | 35 ++++++++ .../blocks/DelimitedParserExecutor.java | 70 +++++++++++++++ .../controllers/BlocksControllerTest.java | 43 +++++++++ .../DelimitedResponseParserTest.java | 43 +++++++++ 7 files changed, 393 insertions(+) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/types/DelimitedParserBlockType.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParserTest.java 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 new file mode 100644 index 0000000..79484e2 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DelimitedParserBlockConfiguration.java @@ -0,0 +1,88 @@ +package it.cnr.isti.workflow.manager.blocks.configurations; + +import java.util.List; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonProperty; + +import it.cnr.isti.workflow.manager.blocks.types.DelimitedParserBlockType; +import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder; +import it.cnr.isti.workflow.manager.configurations.annotations.UiUniqueItemsBy; +import jakarta.validation.constraints.AssertTrue; +import jakarta.validation.Valid; +import jakarta.validation.constraints.Size; +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 DelimitedParserBlockConfiguration extends BlockConfiguration { + + @UiOrder(30) + @Structural + @UiUniqueItemsBy("name") + @Size(max = 10) + @Valid + @JsonProperty(required = true) + private List outputs = List.of(); + + @UiOrder(20) + @Structural + @JsonProperty(required = true) + private String separator; + + @Builder + public DelimitedParserBlockConfiguration(@NonNull String name, String separator, List outputs) { + super(name); + this.separator = separator; + this.outputs = outputs == null ? List.of() : List.copyOf(outputs); + } + + @Override + public Class getBlockType() { + return DelimitedParserBlockType.class; + } + + public static DelimitedParserBlockConfiguration empty() { + DelimitedParserBlockConfiguration configuration = new DelimitedParserBlockConfiguration(); + configuration.name = DelimitedParserBlockType.TYPE; + configuration.separator = null; + configuration.outputs = List.of(new SwitchCase("part1"), new SwitchCase("part2")); + return configuration; + } + + @AssertTrue(message = "separator is required") + @JsonIgnore + boolean isSeparatorValid() { + return separator != null && !separator.isBlank(); + } + + @AssertTrue(message = "outputs are required") + @JsonIgnore + boolean areOutputsValid() { + return outputs != null && !outputs.isEmpty(); + } + + @AssertTrue(message = "outputs must contain unique non-blank names") + @JsonIgnore + boolean areOutputsUnique() { + if (outputs == null || outputs.isEmpty()) { + return true; + } + long distinct = outputs.stream() + .map(SwitchCase::name) + .filter(name -> name != null && !name.isBlank()) + .distinct() + .count(); + long nonBlank = outputs.stream() + .map(SwitchCase::name) + .filter(name -> name != null && !name.isBlank()) + .count(); + return distinct == nonBlank; + } +} \ No newline at end of file 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 new file mode 100644 index 0000000..1b97294 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/DelimitedParserBlockFactory.java @@ -0,0 +1,77 @@ +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.SwitchCase; +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; +import it.cnr.isti.workflow.manager.ios.IOType; + +@Component +public class DelimitedParserBlockFactory implements BlockFactory { + + public static final String INPUT_NAME = "input"; + + private static final List INPUT_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.TEXT, false)); + private static final List OUTPUT_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.TEXT, false), + new IOCapability(IOCapabilityType.BOOLEAN, false), + new IOCapability(IOCapabilityType.ANY, false)); + + @Autowired + private DelimitedParserBlockType blockType; + + @Override + public Block create(DelimitedParserBlockConfiguration configuration) { + Block.BlockBuilder builder = Block.builder() + .input(IODescriptor.input(INPUT_NAME, IOType.TEXT, false, INPUT_CAPABILITIES)) + .specificConfiguration(configuration) + .type(blockType); + + for (SwitchCase 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)))); + } + return builder.build(); + } + + @Override + public Block createEmpty() { + return create(DelimitedParserBlockConfiguration.empty()); + } + + @Override + public Class getBlockType() { + return DelimitedParserBlockType.class; + } + + @Override + public List supportedInputCapabilities() { + return INPUT_CAPABILITIES; + } + + @Override + public List supportedOutputCapabilities() { + return OUTPUT_CAPABILITIES; + } + + private IOCapabilityType toCapabilityType(IOType type) { + return switch (type) { + case BOOLEAN -> IOCapabilityType.BOOLEAN; + case FILE, CSV -> IOCapabilityType.FILE; + case TEXT -> IOCapabilityType.TEXT; + case ANY -> IOCapabilityType.ANY; + }; + } +} \ No newline at end of file diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/DelimitedParserBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/DelimitedParserBlockType.java new file mode 100644 index 0000000..0794c8f --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/DelimitedParserBlockType.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.DelimitedParserBlockConfiguration; + +@Component(DelimitedParserBlockType.TYPE) +public class DelimitedParserBlockType implements BlockType { + + public static final String TYPE = "DelimitedParserBlock"; + + @Override + public String getName() { + return TYPE; + } + + @Override + public String getDescription() { + return "Parses a text payload into configurable typed outputs using a configurable separator"; + } + + @Override + public boolean validate() { + return true; + } + + @Override + public boolean isUserInteractive() { + return false; + } + + @Override + public Class> getBlockConfigurationClass() { + return DelimitedParserBlockConfiguration.class; + } +} \ No newline at end of file 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 new file mode 100644 index 0000000..26d2988 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParser.java @@ -0,0 +1,35 @@ +package it.cnr.isti.workflow.manager.executions.executors; + +import java.util.Arrays; +import java.util.List; +import java.util.regex.Pattern; + +import org.springframework.util.StringUtils; + +public final class DelimitedResponseParser { + + private 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(); + if (parts.size() != expectedParts) { + throw new IllegalArgumentException( + "Expected " + expectedParts + " parts separated by '" + separator + "' but found " + parts.size()); + } + return parts; + } +} \ 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 new file mode 100644 index 0000000..3927ae8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/DelimitedParserExecutor.java @@ -0,0 +1,70 @@ +package it.cnr.isti.workflow.manager.executions.executors.blocks; + +import java.io.File; +import java.util.List; +import java.util.LinkedHashMap; +import java.util.Map; + +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +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.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.executors.DelimitedResponseParser; +import it.cnr.isti.workflow.manager.executions.steps.Input; +import it.cnr.isti.workflow.manager.ios.IOType; + +@Component +public class DelimitedParserExecutor implements BlockExecutor { + + @Override + public Map execute(Block block, List inputs, + Map authorizations, Map executionVariables, + Map executionVariableDescriptors, ExecutionEventLogger eventLogger) { + DelimitedParserBlockConfiguration configuration = (DelimitedParserBlockConfiguration) block.getSpecificConfiguration(); + String payload = inputs.stream() + .filter(input -> 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)); + + List outputCases = configuration.getOutputs() == null ? List.of() : configuration.getOutputs(); + List parts = DelimitedResponseParser.parseParts(payload, configuration.getSeparator(), outputCases.size()); + + Map result = new LinkedHashMap<>(); + for (int index = 0; index < outputCases.size(); index++) { + SwitchCase outputCase = outputCases.get(index); + result.put(outputCase.name(), coerce(parts.get(index), outputCase.type(), outputCase.name())); + } + return result; + } + + @Override + public Class getBlockType() { + return DelimitedParserBlockType.class; + } + + 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"); + } + if (type == IOType.FILE || type == IOType.CSV) { + if (!StringUtils.hasText(rawValue)) { + throw new IllegalArgumentException("Output " + outputName + " expects a non-empty file path"); + } + return new File(rawValue); + } + return rawValue; + } +} \ No newline at end of file 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 df02046..d519e27 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 @@ -32,12 +32,14 @@ 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.controllers.BlocksController.BlockConfigurationDescriptor; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; 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.configurations.annotations.UiContextKeys; import it.cnr.isti.workflow.manager.configurations.retrievers.ExecutionVariablesFieldRetriever; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; @@ -872,6 +874,47 @@ public class BlocksControllerTest { assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("review"))); } + @Test + public void createDelimitedParserBlockExposesTextAndBooleanOutputs() { + Block block = blocksController.create(DelimitedParserBlockConfiguration.builder() + .name("Parser") + .separator("----") + .outputs(List.of( + new SwitchCase("message"), + new SwitchCase("isValid", it.cnr.isti.workflow.manager.ios.IOType.BOOLEAN))) + .build()); + + assertNotNull(block); + assertEquals(DelimitedParserBlockType.TYPE, block.getType().getName()); + assertTrue(block.getInputs().stream().anyMatch(input -> input.getName().equals("input"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("message"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("isValid"))); + assertEquals(it.cnr.isti.workflow.manager.ios.IOType.BOOLEAN, + block.getOutputs().stream() + .filter(output -> output.getName().equals("isValid")) + .findFirst() + .orElseThrow() + .getType()); + } + + @Test + public void delimitedParserBlockSchemaExposesSeparatorAndPropertyOrder() { + BlockConfigurationDescriptor descriptor = blocksController + .getConfigurationDescriptorForType(DelimitedParserBlockType.TYPE); + + assertNotNull(descriptor.schema()); + JsonNode schema = (JsonNode) descriptor.schema(); + + assertEquals( + List.of("name", "separator", "outputs", "type"), + propertyNames(schema.path("properties"))); + assertEquals( + List.of("name", "separator", "outputs", "type"), + arrayValues(schema.path("x-ui-property-order"))); + assertEquals("string", schema.path("properties").path("separator").path("type").asText()); + assertEquals("array", schema.path("properties").path("outputs").path("type").asText()); + } + @Test public void getMCPBridgeExampleForType() { Block block = blocksController.getExampleForType(MCPAgentBlockType.TYPE); diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParserTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParserTest.java new file mode 100644 index 0000000..dc092bf --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/executors/DelimitedResponseParserTest.java @@ -0,0 +1,43 @@ +package it.cnr.isti.workflow.manager.executions.executors; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.util.List; + +import org.junit.jupiter.api.Test; + +class DelimitedResponseParserTest { + + @Test + void parsesTwoPartsWithDefaultLikeSeparator() { + List result = DelimitedResponseParser.parseParts("pippo ---- true", "----", 2); + + assertEquals(List.of("pippo", "true"), result); + } + + @Test + void parsesThreePartsWithCustomSeparator() { + List result = DelimitedResponseParser.parseParts("hello::false::world", "::", 3); + + assertEquals(List.of("hello", "false", "world"), result); + } + + @Test + void rejectsPayloadWithoutSeparator() { + assertThrows(IllegalArgumentException.class, + () -> DelimitedResponseParser.parseParts("hello true", "----", 2)); + } + + @Test + void rejectsWrongNumberOfParts() { + assertThrows(IllegalArgumentException.class, + () -> DelimitedResponseParser.parseParts("hello ---- true", "----", 3)); + } + + @Test + void rejectsInvalidExpectedPartsValue() { + assertThrows(IllegalArgumentException.class, + () -> DelimitedResponseParser.parseParts("hello ---- true", "----", 0)); + } +} \ No newline at end of file From 68c9bca3238b3d8042fcd117a029025e403cc503 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 26 May 2026 23:21:31 +0200 Subject: [PATCH 3/4] Add guard subflow loop feedback mapping --- .../configurations/JsonSchemaProducer.java | 17 +++++-- .../LoopContainerConfiguration.java | 50 ++++++++++++++++++- .../containers/types/LoopContainerType.java | 2 +- .../ContainerSubFlowValidationRegistry.java | 17 ++++--- .../containers/IteratorContainerExecutor.java | 15 +++--- .../containers/LoopContainerExecutor.java | 44 ++++++++++++---- src/main/resources/application.properties | 1 + .../controllers/ContainersControllerTest.java | 12 +++-- .../manager/executions/ExecutionTest.java | 31 ++++++++++++ 9 files changed, 155 insertions(+), 34 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java index 77c478a..cc01bc0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java @@ -55,6 +55,9 @@ import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValida @Component public class JsonSchemaProducer { + private static final String LOOP_BODY_DESCRIPTION = + "Loop body flow. In addition to the base subflow validation, it must expose at least one open non-multiple input that can receive guard feedback. Configure feedbackInput when more than one input is available; it is inferred when there is exactly one."; + private final SchemaGenerator schemaGenerator; private final ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry; @@ -101,7 +104,7 @@ public class JsonSchemaProducer { applyUiRequiredWhenMetadata(root, getMergedMetadata(uiRequiredWhenMap, type)); applyUiUniqueItemsByMetadata(root, getMergedMetadata(uiUniqueItemsByMap, type)); applyUiLabelMetadata(root, getMergedMetadata(uiLabelMap, type)); - applyUiDescriptionMetadata(root, getMergedMetadata(uiDescriptionMap, type)); + applyUiDescriptionMetadata(root, type, getMergedMetadata(uiDescriptionMap, type)); applySizeMetadata(root, getMergedMetadata(sizeMap, type)); applySchemaAllowedValuesMetadata(root, getMergedMetadata(schemaAllowedValuesMap, type)); applyConfigurableAsInputMetadata(root, getMergedMetadata(configurableAsInputMap, type)); @@ -141,7 +144,7 @@ public class JsonSchemaProducer { applyUiRequiredWhenMetadata(classSchema, getMergedMetadata(uiRequiredWhenMap, matchedClass)); applyUiUniqueItemsByMetadata(classSchema, getMergedMetadata(uiUniqueItemsByMap, matchedClass)); applyUiLabelMetadata(classSchema, getMergedMetadata(uiLabelMap, matchedClass)); - applyUiDescriptionMetadata(classSchema, getMergedMetadata(uiDescriptionMap, matchedClass)); + applyUiDescriptionMetadata(classSchema, matchedClass, getMergedMetadata(uiDescriptionMap, matchedClass)); applySizeMetadata(classSchema, getMergedMetadata(sizeMap, matchedClass)); applySchemaAllowedValuesMetadata(classSchema, getMergedMetadata(schemaAllowedValuesMap, matchedClass)); applyConfigurableAsInputMetadata(classSchema, getMergedMetadata(configurableAsInputMap, matchedClass)); @@ -854,7 +857,7 @@ public class JsonSchemaProducer { } } - private void applyUiDescriptionMetadata(ObjectNode classSchema, Map metadata) { + private void applyUiDescriptionMetadata(ObjectNode classSchema, Class ownerClass, Map metadata) { if (metadata == null || metadata.isEmpty()) { return; } @@ -872,6 +875,14 @@ public class JsonSchemaProducer { propertySchema.put("x-ui-description", entry.getValue().value()); } } + containerSubFlowValidationRegistry.typeForConfigurationField(ownerClass, "subFlow") + .filter(ContainerSubFlowValidationType.LOOP_BODY::equals) + .ifPresent(ignored -> { + JsonNode propNode = properties.get("subFlow"); + if (propNode instanceof ObjectNode propertySchema) { + propertySchema.put("x-ui-description", LOOP_BODY_DESCRIPTION); + } + }); } private void applySizeMetadata(ObjectNode classSchema, Map metadata) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java index 287828a..28cd99f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java @@ -1,13 +1,18 @@ package it.cnr.isti.workflow.manager.containers.configurations; +import java.util.List; + import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever; import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen; import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription; import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; +import it.cnr.isti.workflow.manager.configurations.annotations.UiOptionsFromNode; import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; import jakarta.validation.Valid; @@ -26,7 +31,15 @@ public class LoopContainerConfiguration extends ContainerConfiguration 0; } + + @AssertTrue(message = "feedbackInput is required when subFlow exposes more than one open non-multiple input") + @JsonIgnore + boolean isFeedbackInputPresentWhenAmbiguous() { + if (!hasSubFlowNodes() || (feedbackInput != null && !feedbackInput.isBlank())) { + return true; + } + return eligibleFeedbackInputs().size() <= 1; + } + + @AssertTrue(message = "feedbackInput must target an open non-multiple input of the subFlow") + @JsonIgnore + boolean isFeedbackInputValid() { + if (!hasSubFlowNodes()) { + return true; + } + if (feedbackInput == null || feedbackInput.isBlank()) { + return !eligibleFeedbackInputs().isEmpty(); + } + return eligibleFeedbackInputs().stream() + .anyMatch(handle -> handle.publicName().equals(feedbackInput)); + } + + private boolean hasSubFlowNodes() { + return getSubFlow() != null && !getSubFlow().getNodes().isEmpty(); + } + + private List eligibleFeedbackInputs() { + return ContainerFlowInterfaceResolver.getExposedInputs(getSubFlow()).stream() + .filter(handle -> !handle.handle().io().isMultiple()) + .toList(); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java index 24a0120..2242cec 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java @@ -17,7 +17,7 @@ public class LoopContainerType implements ContainerType { @Override public String getDescription() { - return "A container node that repeatedly executes a subflow and then a guard subflow. The guard subflow must expose guard and feedback outputs; guard=true continues the loop with feedback as the next feedback input."; + return "A container node that repeatedly executes a subflow and then a guard subflow. The guard subflow must expose guard and feedback outputs; guard=true continues the loop by passing feedback to the configured internal flow input."; } @Override diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java index 379befa..2d64e75 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java @@ -1,5 +1,6 @@ package it.cnr.isti.workflow.manager.containers.validation; +import java.util.ArrayList; import java.util.EnumMap; import java.util.List; import java.util.Locale; @@ -72,18 +73,18 @@ public class ContainerSubFlowValidationRegistry { List openOutputs) { List exposedOutputs = ContainerFlowInterfaceResolver.getExposedOutputs(subFlow); + List errors = new ArrayList<>(); if (exposedOutputs.isEmpty()) { - return List.of(error("subFlow", + errors.add(error("subFlow", "LoopContainer subFlow must expose at least one open output")); } - boolean hasFeedbackInput = ContainerFlowInterfaceResolver.getExposedInputs(subFlow).stream() - .anyMatch(handle -> LoopContainerConfiguration.FEEDBACK_INPUT.equals(handle.publicName()) - && !handle.handle().io().isMultiple()); - if (hasFeedbackInput) { - return List.of(); + boolean hasFeedbackTarget = ContainerFlowInterfaceResolver.getExposedInputs(subFlow).stream() + .anyMatch(handle -> !handle.handle().io().isMultiple()); + if (!hasFeedbackTarget) { + errors.add(error("subFlow", + "LoopContainer subFlow must expose at least one open non-multiple input that can receive guard feedback")); } - return List.of(error("subFlow", - "LoopContainer subFlow must expose an open non-multiple input named feedback")); + return List.copyOf(errors); } private List validateLoopGuard(FlowData subFlow, diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java index eef0772..d16e429 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -6,6 +6,7 @@ import java.util.List; import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import io.micrometer.core.instrument.MeterRegistry; @@ -30,16 +31,19 @@ import it.cnr.isti.workflow.manager.executions.steps.Input; public class IteratorContainerExecutor implements ContainerExecutor { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(IteratorContainerExecutor.class); - private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private final ExecutionsService executionsService; + private final long innerExecutionWaitIntervalMs; @Autowired(required = false) MeterRegistry meterRegistry; - public IteratorContainerExecutor(ExecutionsService executionsService) { + public IteratorContainerExecutor( + ExecutionsService executionsService, + @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { this.executionsService = executionsService; + this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -143,12 +147,7 @@ public class IteratorContainerExecutor implements ContainerExecutor { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(LoopContainerExecutor.class); - private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private final ExecutionsService executionsService; + private final long innerExecutionWaitIntervalMs; @Autowired(required = false) MeterRegistry meterRegistry; - public LoopContainerExecutor(ExecutionsService executionsService) { + public LoopContainerExecutor( + ExecutionsService executionsService, + @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { this.executionsService = executionsService; + this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -58,6 +62,7 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName = ContainerFlowInterfaceResolver .getExposedInputs(configuration.getSubFlow()).stream() .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed)); + String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); List exposedOutputs = ContainerFlowInterfaceResolver .getExposedOutputs(configuration.getSubFlow()); @@ -112,7 +117,7 @@ public class LoopContainerExecutor implements ContainerExecutor(latestOutputs); } @@ -136,6 +141,32 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName) { + if (configuration.getFeedbackInput() != null && !configuration.getFeedbackInput().isBlank()) { + String configuredInput = configuration.getFeedbackInput(); + ContainerFlowInterfaceResolver.ExposedHandle handle = inputPortsByName.get(configuredInput); + if (handle == null || handle.handle().io().isMultiple()) { + throw new IllegalArgumentException( + "LoopContainer feedbackInput must target an open non-multiple subFlow input: " + configuredInput); + } + return configuredInput; + } + + List eligibleInputs = inputPortsByName.values().stream() + .filter(handle -> !handle.handle().io().isMultiple()) + .toList(); + if (eligibleInputs.size() == 1) { + return eligibleInputs.getFirst().publicName(); + } + if (eligibleInputs.isEmpty()) { + throw new IllegalArgumentException( + "LoopContainer subFlow must expose at least one open non-multiple input to receive guard feedback"); + } + throw new IllegalArgumentException( + "LoopContainer feedbackInput is required when subFlow exposes more than one open non-multiple input"); + } + private GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName, Map iterationInputs, Map latestOutputs, Map authorizations, Map executionVariables, @@ -265,12 +296,7 @@ public class LoopContainerExecutor implements ContainerExecutor "LoopContainerGuard".equals(event.getDetails().get("containerType")))); } + @Test + public void loopCanMapGuardFeedbackToConfiguredInternalInput() { + Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Function Builder") + .llmDescriptor(llmBrick) + .prompt("Build function for request: ${{specification}} in ${{language}}") + .build()); + + Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() + .name("Loop") + .subFlow(FlowData.builder().block(internalBlock).build()) + .guardSubFlow(loopGuardSubFlow()) + .feedbackInput("specification") + .maxIterations(3) + .build()); + + FlowData flow = FlowData.builder().container(container).build(); + ExecutionObject execObject = executionsService.createExecution("Loop mapped feedback flow", flow); + + executionsService.prepareInput(execObject.getId(), container.getId(), "specification", "Implement the same function"); + executionsService.prepareInput(execObject.getId(), container.getId(), "language", "Java"); + execObject = executionsService.startExecution(execObject.getId()); + while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { + execObject = executionsService.getExecution(execObject.getId()); + } + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); + assertEquals("function v2", output); + } + @Test public void singleTextInputRejectsMultipleValues() { Flow flow = flowTestCreator.createFlowwithLLMUnpromptedWithConnection(llmBrick); From 84ff2daa053ca68bf4f6a02a61fc7518ebb3dbf3 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 26 May 2026 23:22:01 +0200 Subject: [PATCH 4/4] Add static favicon --- src/main/resources/static/favicon.ico | Bin 0 -> 1150 bytes 1 file changed, 0 insertions(+), 0 deletions(-) create mode 100644 src/main/resources/static/favicon.ico diff --git a/src/main/resources/static/favicon.ico b/src/main/resources/static/favicon.ico new file mode 100644 index 0000000000000000000000000000000000000000..20a2c13f584c47de33eebaaae051556b1ebba1bc GIT binary patch literal 1150 zcmZQzU<5(|0R|wcz>vYhz#zuJz@P!dKp~(AL>x#lH~{6)ft8!||KWoBZ(RCM494YO bV)TO4jOxdpW=4AW;Yt^SSscAQAe9dQh&Mf1 literal 0 HcmV?d00001