From 68c9bca3238b3d8042fcd117a029025e403cc503 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 26 May 2026 23:21:31 +0200 Subject: [PATCH] Add guard subflow loop feedback mapping --- .../configurations/JsonSchemaProducer.java | 17 +++++-- .../LoopContainerConfiguration.java | 50 ++++++++++++++++++- .../containers/types/LoopContainerType.java | 2 +- .../ContainerSubFlowValidationRegistry.java | 17 ++++--- .../containers/IteratorContainerExecutor.java | 15 +++--- .../containers/LoopContainerExecutor.java | 44 ++++++++++++---- src/main/resources/application.properties | 1 + .../controllers/ContainersControllerTest.java | 12 +++-- .../manager/executions/ExecutionTest.java | 31 ++++++++++++ 9 files changed, 155 insertions(+), 34 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java index 77c478a..cc01bc0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java @@ -55,6 +55,9 @@ import it.cnr.isti.workflow.manager.containers.validation.ContainerSubFlowValida @Component public class JsonSchemaProducer { + private static final String LOOP_BODY_DESCRIPTION = + "Loop body flow. In addition to the base subflow validation, it must expose at least one open non-multiple input that can receive guard feedback. Configure feedbackInput when more than one input is available; it is inferred when there is exactly one."; + private final SchemaGenerator schemaGenerator; private final ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry; @@ -101,7 +104,7 @@ public class JsonSchemaProducer { applyUiRequiredWhenMetadata(root, getMergedMetadata(uiRequiredWhenMap, type)); applyUiUniqueItemsByMetadata(root, getMergedMetadata(uiUniqueItemsByMap, type)); applyUiLabelMetadata(root, getMergedMetadata(uiLabelMap, type)); - applyUiDescriptionMetadata(root, getMergedMetadata(uiDescriptionMap, type)); + applyUiDescriptionMetadata(root, type, getMergedMetadata(uiDescriptionMap, type)); applySizeMetadata(root, getMergedMetadata(sizeMap, type)); applySchemaAllowedValuesMetadata(root, getMergedMetadata(schemaAllowedValuesMap, type)); applyConfigurableAsInputMetadata(root, getMergedMetadata(configurableAsInputMap, type)); @@ -141,7 +144,7 @@ public class JsonSchemaProducer { applyUiRequiredWhenMetadata(classSchema, getMergedMetadata(uiRequiredWhenMap, matchedClass)); applyUiUniqueItemsByMetadata(classSchema, getMergedMetadata(uiUniqueItemsByMap, matchedClass)); applyUiLabelMetadata(classSchema, getMergedMetadata(uiLabelMap, matchedClass)); - applyUiDescriptionMetadata(classSchema, getMergedMetadata(uiDescriptionMap, matchedClass)); + applyUiDescriptionMetadata(classSchema, matchedClass, getMergedMetadata(uiDescriptionMap, matchedClass)); applySizeMetadata(classSchema, getMergedMetadata(sizeMap, matchedClass)); applySchemaAllowedValuesMetadata(classSchema, getMergedMetadata(schemaAllowedValuesMap, matchedClass)); applyConfigurableAsInputMetadata(classSchema, getMergedMetadata(configurableAsInputMap, matchedClass)); @@ -854,7 +857,7 @@ public class JsonSchemaProducer { } } - private void applyUiDescriptionMetadata(ObjectNode classSchema, Map metadata) { + private void applyUiDescriptionMetadata(ObjectNode classSchema, Class ownerClass, Map metadata) { if (metadata == null || metadata.isEmpty()) { return; } @@ -872,6 +875,14 @@ public class JsonSchemaProducer { propertySchema.put("x-ui-description", entry.getValue().value()); } } + containerSubFlowValidationRegistry.typeForConfigurationField(ownerClass, "subFlow") + .filter(ContainerSubFlowValidationType.LOOP_BODY::equals) + .ifPresent(ignored -> { + JsonNode propNode = properties.get("subFlow"); + if (propNode instanceof ObjectNode propertySchema) { + propertySchema.put("x-ui-description", LOOP_BODY_DESCRIPTION); + } + }); } private void applySizeMetadata(ObjectNode classSchema, Map metadata) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java index 287828a..28cd99f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/LoopContainerConfiguration.java @@ -1,13 +1,18 @@ package it.cnr.isti.workflow.manager.containers.configurations; +import java.util.List; + import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever; import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen; import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription; import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; +import it.cnr.isti.workflow.manager.configurations.annotations.UiOptionsFromNode; import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; import jakarta.validation.Valid; @@ -26,7 +31,15 @@ public class LoopContainerConfiguration extends ContainerConfiguration 0; } + + @AssertTrue(message = "feedbackInput is required when subFlow exposes more than one open non-multiple input") + @JsonIgnore + boolean isFeedbackInputPresentWhenAmbiguous() { + if (!hasSubFlowNodes() || (feedbackInput != null && !feedbackInput.isBlank())) { + return true; + } + return eligibleFeedbackInputs().size() <= 1; + } + + @AssertTrue(message = "feedbackInput must target an open non-multiple input of the subFlow") + @JsonIgnore + boolean isFeedbackInputValid() { + if (!hasSubFlowNodes()) { + return true; + } + if (feedbackInput == null || feedbackInput.isBlank()) { + return !eligibleFeedbackInputs().isEmpty(); + } + return eligibleFeedbackInputs().stream() + .anyMatch(handle -> handle.publicName().equals(feedbackInput)); + } + + private boolean hasSubFlowNodes() { + return getSubFlow() != null && !getSubFlow().getNodes().isEmpty(); + } + + private List eligibleFeedbackInputs() { + return ContainerFlowInterfaceResolver.getExposedInputs(getSubFlow()).stream() + .filter(handle -> !handle.handle().io().isMultiple()) + .toList(); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java index 24a0120..2242cec 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/types/LoopContainerType.java @@ -17,7 +17,7 @@ public class LoopContainerType implements ContainerType { @Override public String getDescription() { - return "A container node that repeatedly executes a subflow and then a guard subflow. The guard subflow must expose guard and feedback outputs; guard=true continues the loop with feedback as the next feedback input."; + return "A container node that repeatedly executes a subflow and then a guard subflow. The guard subflow must expose guard and feedback outputs; guard=true continues the loop by passing feedback to the configured internal flow input."; } @Override diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java index 379befa..2d64e75 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/validation/ContainerSubFlowValidationRegistry.java @@ -1,5 +1,6 @@ package it.cnr.isti.workflow.manager.containers.validation; +import java.util.ArrayList; import java.util.EnumMap; import java.util.List; import java.util.Locale; @@ -72,18 +73,18 @@ public class ContainerSubFlowValidationRegistry { List openOutputs) { List exposedOutputs = ContainerFlowInterfaceResolver.getExposedOutputs(subFlow); + List errors = new ArrayList<>(); if (exposedOutputs.isEmpty()) { - return List.of(error("subFlow", + errors.add(error("subFlow", "LoopContainer subFlow must expose at least one open output")); } - boolean hasFeedbackInput = ContainerFlowInterfaceResolver.getExposedInputs(subFlow).stream() - .anyMatch(handle -> LoopContainerConfiguration.FEEDBACK_INPUT.equals(handle.publicName()) - && !handle.handle().io().isMultiple()); - if (hasFeedbackInput) { - return List.of(); + boolean hasFeedbackTarget = ContainerFlowInterfaceResolver.getExposedInputs(subFlow).stream() + .anyMatch(handle -> !handle.handle().io().isMultiple()); + if (!hasFeedbackTarget) { + errors.add(error("subFlow", + "LoopContainer subFlow must expose at least one open non-multiple input that can receive guard feedback")); } - return List.of(error("subFlow", - "LoopContainer subFlow must expose an open non-multiple input named feedback")); + return List.copyOf(errors); } private List validateLoopGuard(FlowData subFlow, diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java index eef0772..d16e429 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -6,6 +6,7 @@ import java.util.List; import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import io.micrometer.core.instrument.MeterRegistry; @@ -30,16 +31,19 @@ import it.cnr.isti.workflow.manager.executions.steps.Input; public class IteratorContainerExecutor implements ContainerExecutor { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(IteratorContainerExecutor.class); - private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private final ExecutionsService executionsService; + private final long innerExecutionWaitIntervalMs; @Autowired(required = false) MeterRegistry meterRegistry; - public IteratorContainerExecutor(ExecutionsService executionsService) { + public IteratorContainerExecutor( + ExecutionsService executionsService, + @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { this.executionsService = executionsService; + this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -143,12 +147,7 @@ public class IteratorContainerExecutor implements ContainerExecutor { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(LoopContainerExecutor.class); - private static final long INNER_EXECUTION_WAIT_TIMEOUT_MS = 60_000L; private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; private final ExecutionsService executionsService; + private final long innerExecutionWaitIntervalMs; @Autowired(required = false) MeterRegistry meterRegistry; - public LoopContainerExecutor(ExecutionsService executionsService) { + public LoopContainerExecutor( + ExecutionsService executionsService, + @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { this.executionsService = executionsService; + this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -58,6 +62,7 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName = ContainerFlowInterfaceResolver .getExposedInputs(configuration.getSubFlow()).stream() .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed)); + String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); List exposedOutputs = ContainerFlowInterfaceResolver .getExposedOutputs(configuration.getSubFlow()); @@ -112,7 +117,7 @@ public class LoopContainerExecutor implements ContainerExecutor(latestOutputs); } @@ -136,6 +141,32 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName) { + if (configuration.getFeedbackInput() != null && !configuration.getFeedbackInput().isBlank()) { + String configuredInput = configuration.getFeedbackInput(); + ContainerFlowInterfaceResolver.ExposedHandle handle = inputPortsByName.get(configuredInput); + if (handle == null || handle.handle().io().isMultiple()) { + throw new IllegalArgumentException( + "LoopContainer feedbackInput must target an open non-multiple subFlow input: " + configuredInput); + } + return configuredInput; + } + + List eligibleInputs = inputPortsByName.values().stream() + .filter(handle -> !handle.handle().io().isMultiple()) + .toList(); + if (eligibleInputs.size() == 1) { + return eligibleInputs.getFirst().publicName(); + } + if (eligibleInputs.isEmpty()) { + throw new IllegalArgumentException( + "LoopContainer subFlow must expose at least one open non-multiple input to receive guard feedback"); + } + throw new IllegalArgumentException( + "LoopContainer feedbackInput is required when subFlow exposes more than one open non-multiple input"); + } + private GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName, Map iterationInputs, Map latestOutputs, Map authorizations, Map executionVariables, @@ -265,12 +296,7 @@ public class LoopContainerExecutor implements ContainerExecutor "LoopContainerGuard".equals(event.getDetails().get("containerType")))); } + @Test + public void loopCanMapGuardFeedbackToConfiguredInternalInput() { + Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Function Builder") + .llmDescriptor(llmBrick) + .prompt("Build function for request: ${{specification}} in ${{language}}") + .build()); + + Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() + .name("Loop") + .subFlow(FlowData.builder().block(internalBlock).build()) + .guardSubFlow(loopGuardSubFlow()) + .feedbackInput("specification") + .maxIterations(3) + .build()); + + FlowData flow = FlowData.builder().container(container).build(); + ExecutionObject execObject = executionsService.createExecution("Loop mapped feedback flow", flow); + + executionsService.prepareInput(execObject.getId(), container.getId(), "specification", "Implement the same function"); + executionsService.prepareInput(execObject.getId(), container.getId(), "language", "Java"); + execObject = executionsService.startExecution(execObject.getId()); + while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { + execObject = executionsService.getExecution(execObject.getId()); + } + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); + assertEquals("function v2", output); + } + @Test public void singleTextInputRejectsMultipleValues() { Flow flow = flowTestCreator.createFlowwithLLMUnpromptedWithConnection(llmBrick);