Add guard subflow loop feedback mapping

This commit is contained in:
Lucio Lelii 2026-05-26 23:21:31 +02:00
parent 83c76ebd4f
commit 68c9bca323
9 changed files with 155 additions and 34 deletions

View File

@ -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<String, UiDescription> metadata) {
private void applyUiDescriptionMetadata(ObjectNode classSchema, Class<?> ownerClass, Map<String, UiDescription> 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<String, Size> metadata) {

View File

@ -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<LoopConta
public static final String GUARD_OUTPUT = "guard";
public static final String FEEDBACK_OUTPUT = "feedback";
public static final String FEEDBACK_INPUT = "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
@ -52,9 +65,11 @@ public class LoopContainerConfiguration extends ContainerConfiguration<LoopConta
public LoopContainerConfiguration(@NonNull String name,
FlowData subFlow,
FlowData guardSubFlow,
String feedbackInput,
Integer maxIterations) {
super(name, subFlow);
this.guardSubFlow = guardSubFlow == null ? FlowData.builder().build() : guardSubFlow;
this.feedbackInput = feedbackInput;
this.maxIterations = maxIterations == null ? 10 : maxIterations;
}
@ -68,6 +83,7 @@ public class LoopContainerConfiguration extends ContainerConfiguration<LoopConta
LoopContainerType.TYPE,
FlowData.builder().build(),
FlowData.builder().build(),
null,
10);
}
@ -76,4 +92,36 @@ public class LoopContainerConfiguration extends ContainerConfiguration<LoopConta
boolean isMaxIterationsValid() {
return maxIterations != null && maxIterations > 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<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 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

View File

@ -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<ContainerFlowInterfaceResolver.OpenHandle> openOutputs) {
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs =
ContainerFlowInterfaceResolver.getExposedOutputs(subFlow);
List<ValidationError> 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<ValidationError> validateLoopGuard(FlowData subFlow,

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) {

View File

@ -6,6 +6,7 @@ import java.util.Map;
import java.util.stream.Collectors;
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;
@ -31,16 +32,19 @@ import it.cnr.isti.workflow.manager.executions.steps.Input;
public class LoopContainerExecutor implements ContainerExecutor<LoopContainerType> {
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<LoopContainerTyp
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());
@ -112,7 +117,7 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
decision.shouldContinue()));
}
if (decision.shouldContinue()) {
currentInputs.put(LoopContainerConfiguration.FEEDBACK_INPUT, decision.feedback());
currentInputs.put(feedbackInput, decision.feedback());
} else {
return new LinkedHashMap<>(latestOutputs);
}
@ -136,6 +141,32 @@ public class LoopContainerExecutor implements ContainerExecutor<LoopContainerTyp
}
}
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;
}
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 GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName,
Map<String, Object> iterationInputs, Map<String, Object> latestOutputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables,
@ -265,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) {

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}

View File

@ -228,18 +228,22 @@ public class ContainersControllerTest {
JsonNode schema = (JsonNode) descriptor.schema();
JsonNode subFlow = schema.path("properties").path("subFlow");
JsonNode feedbackInput = schema.path("properties").path("feedbackInput");
JsonNode guardSubFlow = schema.path("properties").path("guardSubFlow");
JsonNode maxIterations = schema.path("properties").path("maxIterations");
assertEquals(
List.of("name", "subFlow", "guardSubFlow", "maxIterations", "containerType", "type"),
List.of("name", "subFlow", "feedbackInput", "guardSubFlow", "maxIterations", "containerType", "type"),
propertyNames(schema.path("properties")));
assertEquals(
List.of("name", "subFlow", "guardSubFlow", "maxIterations", "containerType", "type"),
List.of("name", "subFlow", "feedbackInput", "guardSubFlow", "maxIterations", "containerType", "type"),
arrayValues(schema.path("x-ui-property-order")));
assertEquals(3, guardSubFlow.path("x-ui-order").asInt());
assertEquals(4, maxIterations.path("x-ui-order").asInt());
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(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());

View File

@ -1454,6 +1454,37 @@ public class ExecutionTest {
.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
public void singleTextInputRejectsMultipleValues() {
Flow flow = flowTestCreator.createFlowwithLLMUnpromptedWithConnection(llmBrick);