Merge branch 'feature/loop-container-guard-subflow'

This commit is contained in:
Lucio Lelii 2026-05-27 10:40:12 +02:00
commit a9ea3d5401
37 changed files with 1130 additions and 515 deletions

View File

@ -2,6 +2,7 @@ package it.cnr.isti.workflow.manager.blocks;
public enum IOCapabilityType {
TEXT,
BOOLEAN,
FILE,
ANY
}

View File

@ -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<DelimitedParserBlockType> {
@UiOrder(30)
@Structural
@UiUniqueItemsBy("name")
@Size(max = 10)
@Valid
@JsonProperty(required = true)
private List<SwitchCase> outputs = List.of();
@UiOrder(20)
@Structural
@JsonProperty(required = true)
private String separator;
@Builder
public DelimitedParserBlockConfiguration(@NonNull String name, String separator, List<SwitchCase> outputs) {
super(name);
this.separator = separator;
this.outputs = outputs == null ? List.of() : List.copyOf(outputs);
}
@Override
public Class<DelimitedParserBlockType> 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;
}
}

View File

@ -49,13 +49,20 @@ 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 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.";
public JsonSchemaProducer() {
private final SchemaGenerator schemaGenerator;
private final ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry;
public JsonSchemaProducer(ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry) {
this.containerSubFlowValidationRegistry = containerSubFlowValidationRegistry;
SchemaGeneratorConfigBuilder configBuilder = new SchemaGeneratorConfigBuilder(
SchemaVersion.DRAFT_7, OptionPreset.PLAIN_JSON)
.with(Option.DEFINITIONS_FOR_ALL_OBJECTS)
@ -97,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));
@ -137,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));
@ -397,9 +404,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<Class<?>, Map<String, LongText>> collectLongTextMetadata(Class<?> rootClass) {
Map<Class<?>, Map<String, LongText>> result = new HashMap<>();
Set<Class<?>> visited = new HashSet<>();
@ -827,7 +857,7 @@ public class JsonSchemaProducer {
}
}
private void applyUiDescriptionMetadata(ObjectNode classSchema, Map<String, UiDescription> metadata) {
private void applyUiDescriptionMetadata(ObjectNode classSchema, Class<?> ownerClass, Map<String, UiDescription> metadata) {
if (metadata == null || metadata.isEmpty()) {
return;
}
@ -845,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<String, Size> metadata) {

View File

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

View File

@ -85,6 +85,9 @@ public interface BlockFactory<T extends BlockType, C extends BlockConfiguration<
if (type == IOType.TEXT) {
return IOCapabilityType.TEXT;
}
if (type == IOType.BOOLEAN) {
return IOCapabilityType.BOOLEAN;
}
if (type == IOType.FILE || type == IOType.CSV) {
return IOCapabilityType.FILE;
}

View File

@ -116,6 +116,7 @@ public class ChatInteractionBlockFactory
return switch (type) {
case FILE, CSV -> IOCapabilityType.FILE;
case TEXT -> IOCapabilityType.TEXT;
case BOOLEAN -> IOCapabilityType.BOOLEAN;
case ANY -> IOCapabilityType.ANY;
};
}

View File

@ -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<DelimitedParserBlockType, DelimitedParserBlockConfiguration> {
public static final String INPUT_NAME = "input";
private static final List<IOCapability> INPUT_CAPABILITIES = List.of(
new IOCapability(IOCapabilityType.TEXT, false));
private static final List<IOCapability> 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<DelimitedParserBlockType> create(DelimitedParserBlockConfiguration configuration) {
Block.BlockBuilder<DelimitedParserBlockType> builder = Block.<DelimitedParserBlockType>builder()
.input(IODescriptor.input(INPUT_NAME, IOType.TEXT, false, INPUT_CAPABILITIES))
.specificConfiguration(configuration)
.type(blockType);
for (SwitchCase outputCase : configuration.getOutputs() == null ? List.<SwitchCase>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<DelimitedParserBlockType> createEmpty() {
return create(DelimitedParserBlockConfiguration.empty());
}
@Override
public Class<DelimitedParserBlockType> getBlockType() {
return DelimitedParserBlockType.class;
}
@Override
public List<IOCapability> supportedInputCapabilities() {
return INPUT_CAPABILITIES;
}
@Override
public List<IOCapability> 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;
};
}
}

View File

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

View File

@ -24,7 +24,9 @@ public class SwitchBlockFactory implements BlockFactory<SwitchBlockType, SwitchB
private static final List<IOCapability> INPUT_CAPABILITIES = List.of(
new IOCapability(IOCapabilityType.TEXT, false),
new IOCapability(IOCapabilityType.TEXT, true));
private static final List<IOCapability> OUTPUT_CAPABILITIES = List.of(new IOCapability(IOCapabilityType.TEXT, false));
private static final List<IOCapability> 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 BlockFactory<SwitchBlockType, SwitchB
.type(blockType);
for (SwitchCase outputCase : configuration.getCases() == null ? List.<SwitchCase>of() : 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<SwitchBlockType, SwitchB
public List<IOCapability> 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;
};
}
}

View File

@ -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<? extends BlockConfiguration<?>> getBlockConfigurationClass() {
return DelimitedParserBlockConfiguration.class;
}
}

View File

@ -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<ValidationError> validationErrors = resolveValidationErrors(flow, context);
private RetrieverItem toItem(FlowEntity flow, String context, String validationType, boolean includeValidation) {
List<ValidationError> validationErrors = resolveValidationErrors(flow, context, validationType);
boolean valid = validationErrors.isEmpty();
List<ValidationError> errorsPayload = includeValidation ? validationErrors : List.of();
@ -78,7 +79,15 @@ public class FlowsFieldRetriever implements SecureDynamicFieldRetriever {
errorsPayload);
}
private List<ValidationError> resolveValidationErrors(FlowEntity flow, String context) {
private List<ValidationError> 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();
}

View File

@ -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<ValidationError> errors = new ArrayList<>(flowExecutionValidator.collectErrors(subFlow));
List<ContainerFlowInterfaceResolver.OpenHandle> openInputs = ContainerFlowInterfaceResolver.getOpenInputs(subFlow);
List<ContainerFlowInterfaceResolver.OpenHandle> 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(),

View File

@ -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<T extends ContainerType> {
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) {

View File

@ -1,19 +1,20 @@
package it.cnr.isti.workflow.manager.containers.configurations;
import com.fasterxml.jackson.annotation.JsonAlias;
import java.util.List;
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.UiOrder;
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.UiRequiredWhen;
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 it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import jakarta.validation.Valid;
import jakarta.validation.constraints.AssertTrue;
import jakarta.validation.constraints.Min;
@ -28,89 +29,47 @@ import lombok.NonNull;
@EqualsAndHashCode(callSuper = true)
public class LoopContainerConfiguration extends ContainerConfiguration<LoopContainerType> {
public static final String GUARD_OUTPUT = "guard";
public static final String FEEDBACK_OUTPUT = "feedback";
@UiOrder(25)
@Structural
@UiEnabledWhen(field = "subFlow", present = true)
@UiOptionsFromNode(collection = "inputs", valueField = "name", labelField = "name")
@JsonProperty(required = false)
@UiLabel("Feedback Input")
@UiDescription("Internal flow input that receives guardSubFlow.feedback on the next loop iteration. Required when the internal flow exposes more than one input; inferred when it exposes exactly one.")
private String feedbackInput;
@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,
FlowData guardSubFlow,
String feedbackInput,
LLMDescriptor feedbackLlmDescriptor,
String feedbackPrompt,
Integer maxIterations) {
super(name, subFlow);
this.guardCondition = guardCondition;
this.useLlm = useLlm;
this.llmDescriptor = llmDescriptor;
this.guardPrompt = guardPrompt;
this.guardSubFlow = guardSubFlow == null ? FlowData.builder().build() : guardSubFlow;
this.feedbackInput = feedbackInput;
this.feedbackLlmDescriptor = feedbackLlmDescriptor;
this.feedbackPrompt = feedbackPrompt;
this.maxIterations = maxIterations == null ? 10 : maxIterations;
}
@ -123,89 +82,46 @@ public class LoopContainerConfiguration extends ContainerConfiguration<LoopConta
return new LoopContainerConfiguration(
LoopContainerType.TYPE,
FlowData.builder().build(),
"",
false,
null,
"",
null,
null,
FlowData.builder().build(),
null,
10);
}
@AssertTrue(message = "guardCondition is required when useLlm is false")
@JsonIgnore
boolean isDeterministicConfigurationValid() {
return useLlm || (guardCondition != null && !guardCondition.isBlank());
}
@AssertTrue(message = "llmDescriptor is required when useLlm is true")
@JsonIgnore
boolean isLlmConfigurationValid() {
return !useLlm || llmDescriptor != null;
}
@AssertTrue(message = "guardPrompt is required when useLlm is true")
@JsonIgnore
boolean isLlmPromptConfigurationValid() {
return !useLlm || (guardPrompt != null && !guardPrompt.isBlank());
}
@AssertTrue(message = "feedbackPrompt requires feedbackInput")
@JsonIgnore
boolean isFeedbackInputPresentWhenFeedbackPromptConfigured() {
return feedbackPrompt == null || feedbackPrompt.isBlank()
|| (feedbackInput != null && !feedbackInput.isBlank());
}
@AssertTrue(message = "feedbackInput requires feedbackPrompt")
@JsonIgnore
boolean isFeedbackPromptPresentWhenFeedbackInputConfigured() {
return feedbackInput == null || feedbackInput.isBlank()
|| (feedbackPrompt != null && !feedbackPrompt.isBlank());
}
@AssertTrue(message = "feedbackLlmDescriptor is required when feedbackPrompt is configured")
@JsonIgnore
boolean isFeedbackLlmConfigurationValid() {
return feedbackPrompt == null || feedbackPrompt.isBlank() || feedbackLlmDescriptor != null;
}
@AssertTrue(message = "feedbackInput must target an open non-multiple input of the subFlow")
@JsonIgnore
boolean isFeedbackInputValid() {
if (!hasSubFlowNodes() || feedbackInput == null || feedbackInput.isBlank()) {
return true;
}
return ContainerFlowInterfaceResolver.getExposedInputs(getSubFlow()).stream()
.anyMatch(handle -> 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")
@AssertTrue(message = "feedbackInput is required when subFlow exposes more than one open non-multiple input")
@JsonIgnore
boolean hasOpenOutputsWhenConfigured() {
return !hasSubFlowNodes() || !ContainerFlowInterfaceResolver.getExposedOutputs(getSubFlow()).isEmpty();
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();
}
@Deprecated
@JsonIgnore
public String getCondition() {
return guardCondition;
}
@Deprecated
@JsonIgnore
public String getPrompt() {
return guardPrompt;
private List<ContainerFlowInterfaceResolver.ExposedHandle> eligibleFeedbackInputs() {
return ContainerFlowInterfaceResolver.getExposedInputs(getSubFlow()).stream()
.filter(handle -> !handle.handle().io().isMultiple())
.toList();
}
}

View File

@ -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 by passing feedback to the configured internal flow input.";
}
@Override

View File

@ -0,0 +1,122 @@
package it.cnr.isti.workflow.manager.containers.validation;
import java.util.ArrayList;
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<ContainerSubFlowValidationType, Rule> 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<ContainerSubFlowValidationType> 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<ValidationError> validate(ContainerSubFlowValidationType type, FlowData subFlow,
List<ContainerFlowInterfaceResolver.OpenHandle> openInputs,
List<ContainerFlowInterfaceResolver.OpenHandle> openOutputs) {
Rule rule = rules.get(type == null ? ContainerSubFlowValidationType.CONTAINER : type);
return rule == null ? List.of() : rule.validate(subFlow, openInputs, openOutputs);
}
private List<ValidationError> validateLoopBody(FlowData subFlow,
List<ContainerFlowInterfaceResolver.OpenHandle> openInputs,
List<ContainerFlowInterfaceResolver.OpenHandle> openOutputs) {
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs =
ContainerFlowInterfaceResolver.getExposedOutputs(subFlow);
List<ValidationError> errors = new ArrayList<>();
if (exposedOutputs.isEmpty()) {
errors.add(error("subFlow",
"LoopContainer subFlow must expose at least one open output"));
}
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.copyOf(errors);
}
private List<ValidationError> validateLoopGuard(FlowData subFlow,
List<ContainerFlowInterfaceResolver.OpenHandle> openInputs,
List<ContainerFlowInterfaceResolver.OpenHandle> openOutputs) {
List<ContainerFlowInterfaceResolver.ExposedHandle> 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<ContainerFlowInterfaceResolver.ExposedHandle> 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<ValidationError> validate(FlowData subFlow,
List<ContainerFlowInterfaceResolver.OpenHandle> openInputs,
List<ContainerFlowInterfaceResolver.OpenHandle> openOutputs);
}
}

View File

@ -0,0 +1,7 @@
package it.cnr.isti.workflow.manager.containers.validation;
public enum ContainerSubFlowValidationType {
CONTAINER,
LOOP_BODY,
LOOP_GUARD
}

View File

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

View File

@ -437,6 +437,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;
};

View File

@ -3,6 +3,7 @@ package it.cnr.isti.workflow.manager.executions;
public enum ExecutionVariableKind {
ANY,
TEXT,
BOOLEAN,
JSON,
FILE_PATH,
HTTP_RESOURCE,

View File

@ -493,9 +493,11 @@ public class ExecutionsService {
collectHttpRequirement(requirements, block);
}
for (Container<?> container : flow.getContainers() == null ? List.<Container<?>>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());
}
}
}
}
@ -564,28 +566,6 @@ public class ExecutionsService {
.addStepReference(block.getId(), block.getName());
}
private void collectContainerRequirement(Map<String, RequirementAccumulator> 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) {

View File

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

View File

@ -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<DelimitedParserBlockType> {
@Override
public Map<String, Object> execute(Block<DelimitedParserBlockType> block, List<Input> inputs,
Map<String, Object> authorizations, Map<String, Object> executionVariables,
Map<String, ExecutionVariableDescriptor> 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<SwitchCase> outputCases = configuration.getOutputs() == null ? List.of() : configuration.getOutputs();
List<String> parts = DelimitedResponseParser.parseParts(payload, configuration.getSeparator(), outputCases.size());
Map<String, Object> 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<DelimitedParserBlockType> 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;
}
}

View File

@ -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<SwitchBlockType> {
? 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<SwitchBlockType> {
return validateSelectedOutput(result.toString(), allowedOutputs);
}
private Object coercePayload(String payload, String selectedOutput, List<IODescriptor> 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<String, Object> inputValues,
Map<String, Object> authorizations, Map<String, Object> executionVariables, Set<String> allowedOutputs) {
LLMDescriptor llmDescriptor = config.getLlmDescriptor();

View File

@ -174,6 +174,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
it.cnr.isti.workflow.manager.ios.IODescriptor descriptor) {
return switch (descriptor.getType()) {
case TEXT -> 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;
};

View File

@ -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<IteratorContainerType> {
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<IteratorCont
innerExecution = executionsService.startExecution(innerExecution.getId());
try {
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
finalStatus = "timeout";
throw new IllegalStateException("IteratorContainer subflow timed out while waiting for completion");
}
innerExecution.getContext().awaitStatusChangeWhileRunning(innerExecutionWaitIntervalMs);
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
@ -238,6 +237,7 @@ public class IteratorContainerExecutor implements ContainerExecutor<IteratorCont
it.cnr.isti.workflow.manager.ios.IODescriptor descriptor) {
return switch (descriptor.getType()) {
case TEXT -> 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;
};

View File

@ -3,18 +3,11 @@ 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.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.Timer;
@ -29,45 +22,29 @@ 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<LoopContainerType> {
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<String, LLMProvider> llmProviders;
private final ExpressionParser expressionParser = new SpelExpressionParser();
private final long innerExecutionWaitIntervalMs;
@Autowired(required = false)
MeterRegistry meterRegistry;
public LoopContainerExecutor(ExecutionsService executionsService, Map<String, LLMProvider> llmProviders) {
public LoopContainerExecutor(
ExecutionsService executionsService,
@Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) {
this.executionsService = executionsService;
this.llmProviders = llmProviders;
this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs;
}
@Override
@ -78,10 +55,14 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
if (configuration.getSubFlow() == null || configuration.getSubFlow().getNodes().isEmpty()) {
throw new IllegalArgumentException("LoopContainer subFlow is empty");
}
if (configuration.getGuardSubFlow() == null || configuration.getGuardSubFlow().getNodes().isEmpty()) {
throw new IllegalArgumentException("LoopContainer guardSubFlow is empty");
}
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
.getExposedInputs(configuration.getSubFlow()).stream()
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed));
String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName);
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs = ContainerFlowInterfaceResolver
.getExposedOutputs(configuration.getSubFlow());
@ -97,39 +78,18 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
"Starting loop iteration " + iteration,
Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration));
}
ExecutionObject innerExecution = executionsService.createExecution(
ExecutionObject innerExecution = executeSubFlow(
container.getName() + " iteration " + iteration,
configuration.getSubFlow());
forwardInnerEvents(innerExecution, eventLogger, iteration, LoopContainerType.TYPE);
propagateAuthorizations(innerExecution, authorizations);
executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors);
Map<String, Object> 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<String, Object> 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 +103,24 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
Map<String, Object> 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(feedbackInput, decision.feedback());
} else {
return new LinkedHashMap<>(latestOutputs);
}
applyFeedbackInputForNextIteration(configuration, currentInputs, iterationInputs, latestOutputs,
authorizations, executionVariables, iteration, eventLogger);
}
throw new IllegalStateException(
@ -178,121 +141,58 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
}
}
private boolean shouldStop(LoopContainerConfiguration configuration, Map<String, Object> currentInputs,
Map<String, Object> latestOutputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, int iteration) {
Map<String, Object> 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<String, Object> values,
Map<String, Object> currentInputs, Map<String, Object> latestOutputs,
Map<String, Object> 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");
private String resolveFeedbackInput(LoopContainerConfiguration configuration,
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> 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;
}
return result;
}
private boolean evaluateWithLlm(LoopContainerConfiguration configuration, Map<String, Object> values,
Map<String, Object> currentInputs, Map<String, Object> latestOutputs,
Map<String, Object> authorizations, Map<String, Object> 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;
List<ContainerFlowInterfaceResolver.ExposedHandle> 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 void applyFeedbackInputForNextIteration(LoopContainerConfiguration configuration, Map<String, Object> nextInputs,
private GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName,
Map<String, Object> iterationInputs, Map<String, Object> latestOutputs, Map<String, Object> authorizations,
Map<String, Object> 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<String, Object> 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<String, Object> values,
Map<String, Object> currentInputs, Map<String, Object> latestOutputs,
Map<String, Object> 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<String, Object> values,
Map<String, Object> iterationInputs, Map<String, Object> latestOutputs,
Map<String, Object> 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<String, Object> executionVariables,
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger,
int iteration) {
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> guardInputsByName = ContainerFlowInterfaceResolver
.getExposedInputs(configuration.getGuardSubFlow()).stream()
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed));
List<ContainerFlowInterfaceResolver.ExposedHandle> guardOutputs = ContainerFlowInterfaceResolver
.getExposedOutputs(configuration.getGuardSubFlow());
Map<String, Object> guardInputValues = buildGuardTemplateValues(iterationInputs, latestOutputs, iteration);
ExecutionObject guardExecution = executeSubFlow(
containerName + " guard iteration " + iteration,
configuration.getGuardSubFlow(),
guardInputValues,
guardInputsByName,
"LoopContainerGuard",
authorizations,
executionVariables,
executionVariableDescriptors,
eventLogger,
iteration);
Map<String, Object> 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<String, Object> buildGuardTemplateValues(Map<String, Object> currentInputs,
@ -310,45 +210,73 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
return values;
}
private LLMProvider resolveProvider(String providerName) {
LLMProvider provider = llmProviders.get(providerName);
if (provider != null) {
return provider;
}
return llmProviders.values().stream()
.filter(candidate -> Objects.equals(candidate.getName(), providerName))
.findFirst()
.orElseThrow(() -> new IllegalArgumentException("Provider not found: " + providerName));
}
private void ensureAuthorization(LLMProvider llmProvider, String providerName, Map<String, Object> 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<String, Object> 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<String, Object> inputValues,
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName,
String containerType,
Map<String, Object> authorizations,
Map<String, Object> executionVariables,
Map<String, ExecutionVariableDescriptor> 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<String, Object> 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<String, ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> collectExposedOutputs(ExecutionObject execution,
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs) {
Map<String, Object> 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<String, Object> 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<String, Object> authorizations) {
@ -368,12 +296,7 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
innerExecution = executionsService.startExecution(innerExecution.getId());
try {
while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) {
ExecutionStatus observedStatus = innerExecution.getContext()
.awaitStatusChangeWhileRunning(INNER_EXECUTION_WAIT_TIMEOUT_MS);
if (observedStatus == ExecutionStatus.RUNNING) {
finalStatus = "timeout";
throw new IllegalStateException("LoopContainer subflow timed out while waiting for completion");
}
innerExecution.getContext().awaitStatusChangeWhileRunning(innerExecutionWaitIntervalMs);
}
if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) {
@ -431,6 +354,7 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
it.cnr.isti.workflow.manager.ios.IODescriptor descriptor) {
return switch (descriptor.getType()) {
case TEXT -> 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 +386,7 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
}
return details;
}
private record GuardDecision(boolean shouldContinue, String feedback) {
}
}

View File

@ -109,6 +109,11 @@ public class Input {
throw new IllegalArgumentException("Input " + descriptor.getName() + " expects a text value");
}
}
case BOOLEAN -> {
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");

View File

@ -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<ValidFlowStructure
@Autowired
List<ContainerFactory<?, ?>> containerFactories;
@Autowired
ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry;
@Override
public boolean isValid(FlowData flowData, ConstraintValidatorContext context) {
if (flowData == null) {
@ -139,24 +146,9 @@ public class FlowDataValidator implements ConstraintValidator<ValidFlowStructure
}
ContainerConfiguration<?> containerConfiguration = container.getSpecificConfiguration();
FlowData subFlow = containerConfiguration.getSubFlow();
if (subFlow != null && !subFlow.getNodes().isEmpty()) {
List<FlowNode> 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<ValidFlowStructure
}
}
private void validateContainerSubFlow(Container<?> container, ContainerConfiguration<?> containerConfiguration,
String fieldName, FlowData subFlow) {
String fieldPath = "specificConfiguration." + fieldName;
if (subFlow == null || subFlow.getNodes().isEmpty()) {
return;
}
List<FlowNode> 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<ValidationError> 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<String, FlowNode> nodesById) {
if (connection == null) {
throw validationError(error(ValidationErrorCode.NULL_CONNECTION, "flow", null, "connections", "Flow contains a null connection"));

View File

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

View File

@ -5,6 +5,7 @@ import com.fasterxml.jackson.annotation.JsonValue;
public enum IOType {
TEXT,
BOOLEAN,
FILE,
CSV,
ANY;

View File

@ -51,6 +51,7 @@ app.mcp.servers.file=${MCP_SERVERS_FILE:}
app.executions.cache.max-size=${APP_EXECUTIONS_CACHE_MAX_SIZE:1000}
app.executions.cache.final-ttl-ms=${APP_EXECUTIONS_CACHE_FINAL_TTL_MS:1800000}
app.executions.cache.cleanup-interval-ms=${APP_EXECUTIONS_CACHE_CLEANUP_INTERVAL_MS:60000}
app.container.subflow.wait-interval-ms=${CONTAINER_SUBFLOW_WAIT_INTERVAL_MS:5000}
app.http.outbound.timeout-seconds=${HTTP_OUTBOUND_TIMEOUT_SECONDS:30}
app.http.outbound.retry-attempts=${HTTP_OUTBOUND_RETRY_ATTEMPTS:1}
app.http.outbound.allow-private-network=${HTTP_OUTBOUND_ALLOW_PRIVATE_NETWORK:false}

Binary file not shown.

After

Width:  |  Height:  |  Size: 1.1 KiB

View File

@ -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<DelimitedParserBlockType> 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<MCPAgentBlockType> block = blocksController.getExampleForType(MCPAgentBlockType.TYPE);

View File

@ -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<SwitchBlockType> 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<SwitchBlockType> 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<ContainersController.ContainerConfigurationDescriptor> 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,38 @@ 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 subFlow = schema.path("properties").path("subFlow");
JsonNode feedbackInput = schema.path("properties").path("feedbackInput");
JsonNode feedbackPrompt = schema.path("properties").path("feedbackPrompt");
JsonNode feedbackLlmDescriptor = schema.path("properties").path("feedbackLlmDescriptor");
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", "feedbackInput", "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", "feedbackInput", "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(3, feedbackInput.path("x-ui-order").asInt());
assertEquals(4, guardSubFlow.path("x-ui-order").asInt());
assertEquals(5, maxIterations.path("x-ui-order").asInt());
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(ContainerSubFlowValidationType.LOOP_BODY.name(), subFlow.path("x-subflow-validation-type").asText());
assertTrue(subFlow.path("x-ui-description").asText().contains("Configure feedbackInput when more than one input is available"));
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 +266,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 +396,29 @@ public class ContainersControllerTest {
.anyMatch(error -> error.message().contains("Interactive nodes inside containers are not supported yet")));
}
@Test
public void validateLoopGuardSubFlowRequiresGuardAndFeedbackOutputs() {
Block<LLMBlockType> 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<GenericContainerType> container = containersController.create(GenericContainerConfiguration.builder()
@ -510,19 +527,19 @@ public class ContainersControllerTest {
.provider("testProvider")
.model("testModel")
.build())
.prompt("Analyze ${{candidate}}")
.prompt("Analyze ${{candidate}} with ${{feedback}}")
.build());
Container<LoopContainerType> 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")));
}

View File

@ -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<LLMBlockType> guardEvaluator = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Guard Evaluator")
.llmDescriptor(llmBrick)
.prompt("Loop guard subflow. Previous output: ${{outputs.response}}")
.build());
Block<SwitchBlockType> 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<LLMBlockType> feedbackBuilder = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Feedback Builder")
.llmDescriptor(llmBrick)
.prompt("Loop feedback subflow.")
.build());
Block<SwitchBlockType> 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<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Internal LLM")
.llmDescriptor(llmBrick)
.prompt("Hello, ${{name}}!")
.prompt("Hello, ${{feedback}}!")
.build());
Container<LoopContainerType> 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<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Internal LLM")
.llmDescriptor(llmBrick)
.prompt("Hello, ${{name}}!")
.prompt("Hello, ${{feedback}}!")
.build());
Container<LoopContainerType> 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<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Internal LLM")
.llmDescriptor(llmBrick)
.prompt("Hello, ${{name}}!")
.prompt("Hello, ${{feedback}}!")
.build());
Container<LoopContainerType> 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<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Function Builder")
.llmDescriptor(llmBrick)
.prompt("Build function for request: ${{request}}")
.prompt("Build function for request: ${{feedback}}")
.build());
Container<LoopContainerType> 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,38 @@ 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
public void loopCanMapGuardFeedbackToConfiguredInternalInput() {
Block<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Function Builder")
.llmDescriptor(llmBrick)
.prompt("Build function for request: ${{specification}} in ${{language}}")
.build());
Container<LoopContainerType> 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

View File

@ -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<String> result = DelimitedResponseParser.parseParts("pippo ---- true", "----", 2);
assertEquals(List.of("pippo", "true"), result);
}
@Test
void parsesThreePartsWithCustomSeparator() {
List<String> 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));
}
}