From 56424d2f0f00beff0e74e54ef2f67b06e6da92ab Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Tue, 17 Mar 2026 10:15:45 +0100 Subject: [PATCH] IteratorContainer added e small packaging refactoring --- .../containers/ContainerSubFlowValidator.java | 1 + .../IteratorContainerConfiguration.java | 54 ++++++ .../factories/GenericContainerFactory.java | 4 +- .../factories/IteratorContainerFactory.java | 52 ++++++ .../ContainerFlowInterface.java | 2 +- .../ContainerFlowInterfaceResolver.java | 15 +- .../IteratorContainerInterfaceResolver.java | 169 ++++++++++++++++++ .../types/IteratorContainerType.java | 37 ++++ .../controllers/ContainersController.java | 24 +-- .../manager/executions/ExecutionsService.java | 4 +- .../containers/GenericContainerExecutor.java | 19 +- .../containers/IteratorContainerExecutor.java | 124 +++++++++++++ .../flows/validation/FlowDataValidator.java | 39 ++-- .../validation/FlowExecutionValidator.java | 13 +- .../controllers/ContainersControllerTest.java | 62 +++++++ .../manager/executions/ExecutionTest.java | 44 ++++- 16 files changed, 622 insertions(+), 41 deletions(-) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/configurations/IteratorContainerConfiguration.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/factories/IteratorContainerFactory.java rename src/main/java/it/cnr/isti/workflow/manager/containers/{ => iresolvers}/ContainerFlowInterface.java (76%) rename src/main/java/it/cnr/isti/workflow/manager/containers/{ => iresolvers}/ContainerFlowInterfaceResolver.java (87%) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/IteratorContainerInterfaceResolver.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/containers/types/IteratorContainerType.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java b/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java index 85885eb..2cf64df 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerSubFlowValidator.java @@ -5,6 +5,7 @@ import java.util.List; import org.springframework.stereotype.Component; +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.FlowExecutionValidator; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/IteratorContainerConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/IteratorContainerConfiguration.java new file mode 100644 index 0000000..3653fad --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/configurations/IteratorContainerConfiguration.java @@ -0,0 +1,54 @@ +package it.cnr.isti.workflow.manager.containers.configurations; + +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.containers.types.IteratorContainerType; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import jakarta.validation.Valid; +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 IteratorContainerConfiguration extends ContainerConfiguration { + + @Structural + @Valid + @FieldRetriever( + name = "Flows", + url = "/secure-retriever/Flows/subFlow/items", + structuredData = true, + requiresAuth = true, + validationUrl = "/containers/validate-subflow") + @JsonProperty(required = false) + private FlowData subFlow; + + @Structural + @JsonProperty(required = false) + private String iterationInput; + + @Builder + public IteratorContainerConfiguration(@NonNull String name, FlowData subFlow, String iterationInput) { + super(name); + this.subFlow = subFlow == null ? FlowData.builder().build() : subFlow; + this.iterationInput = iterationInput; + } + + @Override + public Class getContainerType() { + return IteratorContainerType.class; + } + + public static IteratorContainerConfiguration empty() { + return new IteratorContainerConfiguration( + IteratorContainerType.TYPE, + FlowData.builder().build(), + null); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/factories/GenericContainerFactory.java b/src/main/java/it/cnr/isti/workflow/manager/containers/factories/GenericContainerFactory.java index 2a3235b..ccb7665 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/factories/GenericContainerFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/factories/GenericContainerFactory.java @@ -5,9 +5,9 @@ import org.springframework.stereotype.Component; import it.cnr.isti.workflow.manager.blocks.Position; import it.cnr.isti.workflow.manager.containers.Container; -import it.cnr.isti.workflow.manager.containers.ContainerFlowInterface; -import it.cnr.isti.workflow.manager.containers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterface; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; @Component diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/factories/IteratorContainerFactory.java b/src/main/java/it/cnr/isti/workflow/manager/containers/factories/IteratorContainerFactory.java new file mode 100644 index 0000000..d3f5267 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/factories/IteratorContainerFactory.java @@ -0,0 +1,52 @@ +package it.cnr.isti.workflow.manager.containers.factories; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.Position; +import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterface; +import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; + +@Component +public class IteratorContainerFactory implements ContainerFactory { + + private static final Position DEFAULT_POSITION = new Position(0, 0); + + private final IteratorContainerType containerType; + + public IteratorContainerFactory(IteratorContainerType containerType) { + this.containerType = containerType; + } + + @Override + public Container create(IteratorContainerConfiguration configuration) { + ContainerFlowInterface exposedInterface = isEmpty(configuration) + ? new ContainerFlowInterface(java.util.List.of(), java.util.List.of()) + : IteratorContainerInterfaceResolver.resolve(configuration); + return Container.builder() + .inputs(exposedInterface.inputs()) + .outputs(exposedInterface.outputs()) + .specificConfiguration(configuration) + .type(containerType) + .position(DEFAULT_POSITION) + .build(); + } + + @Override + public Container createEmpty() { + return create(IteratorContainerConfiguration.empty()); + } + + @Override + public Class getContainerType() { + return IteratorContainerType.class; + } + + private boolean isEmpty(IteratorContainerConfiguration configuration) { + return configuration == null + || configuration.getSubFlow() == null + || configuration.getSubFlow().getNodes().isEmpty(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterface.java b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterface.java similarity index 76% rename from src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterface.java rename to src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterface.java index d4ea332..6f29079 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterface.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterface.java @@ -1,4 +1,4 @@ -package it.cnr.isti.workflow.manager.containers; +package it.cnr.isti.workflow.manager.containers.iresolvers; import java.util.List; diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterfaceResolver.java b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java similarity index 87% rename from src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterfaceResolver.java rename to src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java index 4b4dd6d..fb74e0e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/ContainerFlowInterfaceResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java @@ -1,4 +1,4 @@ -package it.cnr.isti.workflow.manager.containers; +package it.cnr.isti.workflow.manager.containers.iresolvers; import java.util.ArrayList; import java.util.LinkedHashMap; @@ -7,7 +7,9 @@ import java.util.Map; import java.util.Optional; import java.util.stream.Collectors; +import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.model.FlowNode; @@ -18,6 +20,17 @@ public final class ContainerFlowInterfaceResolver { private ContainerFlowInterfaceResolver() { } + public static ContainerFlowInterface resolve(ContainerConfiguration configuration) { + if (configuration instanceof GenericContainerConfiguration genericContainerConfiguration) { + return resolve(genericContainerConfiguration); + } + if (configuration instanceof IteratorContainerConfiguration iteratorContainerConfiguration) { + return IteratorContainerInterfaceResolver.resolve(iteratorContainerConfiguration); + } + throw new IllegalArgumentException("Unsupported container configuration: " + + (configuration == null ? "null" : configuration.getClass().getSimpleName())); + } + public static ContainerFlowInterface resolve(GenericContainerConfiguration configuration) { FlowData subFlow = configuration == null ? null : configuration.getSubFlow(); List openInputs = getExposedInputs(subFlow); diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/IteratorContainerInterfaceResolver.java b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/IteratorContainerInterfaceResolver.java new file mode 100644 index 0000000..470921c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/IteratorContainerInterfaceResolver.java @@ -0,0 +1,169 @@ +package it.cnr.isti.workflow.manager.containers.iresolvers; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; + +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; +import it.cnr.isti.workflow.manager.ios.IODescriptor; + +public final class IteratorContainerInterfaceResolver { + + private IteratorContainerInterfaceResolver() { + } + + public static ContainerFlowInterface resolve(IteratorContainerConfiguration configuration) { + Resolution resolution = resolvePorts(configuration); + return new ContainerFlowInterface(resolution.inputs(), resolution.outputs()); + } + + public static Resolution resolvePorts(IteratorContainerConfiguration configuration) { + if (configuration == null) { + throw new IllegalArgumentException("IteratorContainer configuration is required"); + } + + List exposedInputs = ContainerFlowInterfaceResolver + .getExposedInputs(configuration.getSubFlow()); + List exposedOutputs = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getSubFlow()); + + if (exposedInputs.isEmpty()) { + throw new IllegalArgumentException("IteratorContainer subFlow must expose at least one open input"); + } + String iterationInput = configuration.getIterationInput(); + if (iterationInput == null || iterationInput.isBlank()) { + throw new IllegalArgumentException("IteratorContainer iterationInput is required"); + } + + ContainerFlowInterfaceResolver.ExposedHandle iteratedHandle = exposedInputs.stream() + .filter(handle -> handle.publicName().equals(iterationInput)) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException( + "IteratorContainer iterationInput does not match an open subFlow input: " + iterationInput)); + + if (iteratedHandle.handle().io().isMultiple()) { + throw new IllegalArgumentException( + "IteratorContainer iterationInput must target a non-multiple subFlow input: " + iterationInput); + } + + String iteratedPublicName = uniquePluralizedName(iteratedHandle.publicName(), + exposedInputs.stream() + .map(ContainerFlowInterfaceResolver.ExposedHandle::publicName) + .filter(name -> !name.equals(iteratedHandle.publicName())) + .toList()); + + List resolvedInputs = new ArrayList<>(); + for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedInputs) { + if (exposedHandle.equals(iteratedHandle)) { + resolvedInputs.add(new ResolvedInput( + iteratedPublicName, + cloneDescriptor(iteratedPublicName, exposedHandle.handle().io(), true), + exposedHandle, + true)); + } else { + resolvedInputs.add(new ResolvedInput( + exposedHandle.publicName(), + cloneDescriptor(exposedHandle.publicName(), exposedHandle.handle().io(), + exposedHandle.handle().io().isMultiple()), + exposedHandle, + false)); + } + } + + List resolvedOutputs = exposedOutputs.stream() + .map(exposedHandle -> new ResolvedOutput( + exposedHandle.publicName(), + cloneDescriptor(exposedHandle.publicName(), exposedHandle.handle().io(), true), + exposedHandle)) + .toList(); + + return new Resolution( + List.copyOf(resolvedInputs), + List.copyOf(resolvedOutputs), + iteratedHandle, + iteratedPublicName); + } + + private static IODescriptor cloneDescriptor(String name, IODescriptor source, boolean multiple) { + return new IODescriptor(name, source.getType(), multiple, source.getValueKinds()); + } + + private static String uniquePluralizedName(String name, List reservedNames) { + Set reserved = new LinkedHashSet<>(reservedNames); + String candidate = pluralize(name); + if (!reserved.contains(candidate)) { + return candidate; + } + candidate = name + "List"; + if (!reserved.contains(candidate)) { + return candidate; + } + int counter = 2; + while (reserved.contains(candidate + counter)) { + counter++; + } + return candidate + counter; + } + + private static String pluralize(String name) { + int separatorIndex = name.lastIndexOf('.'); + if (separatorIndex >= 0) { + return name.substring(0, separatorIndex + 1) + pluralizeSegment(name.substring(separatorIndex + 1)); + } + return pluralizeSegment(name); + } + + private static String pluralizeSegment(String segment) { + if (segment.endsWith("y") && segment.length() > 1 && !isVowel(segment.charAt(segment.length() - 2))) { + return segment.substring(0, segment.length() - 1) + "ies"; + } + if (segment.endsWith("s") || segment.endsWith("x") || segment.endsWith("z") + || segment.endsWith("ch") || segment.endsWith("sh")) { + return segment + "es"; + } + return segment + "s"; + } + + private static boolean isVowel(char character) { + return switch (Character.toLowerCase(character)) { + case 'a', 'e', 'i', 'o', 'u' -> true; + default -> false; + }; + } + + public record Resolution( + List resolvedInputs, + List resolvedOutputs, + ContainerFlowInterfaceResolver.ExposedHandle iteratedHandle, + String iteratedPublicName) { + + public List inputs() { + return resolvedInputs.stream().map(ResolvedInput::descriptor).toList(); + } + + public List outputs() { + return resolvedOutputs.stream().map(ResolvedOutput::descriptor).toList(); + } + + public Map inputsByPublicName() { + return resolvedInputs.stream() + .collect(LinkedHashMap::new, (map, input) -> map.put(input.publicName(), input), Map::putAll); + } + } + + public record ResolvedInput( + String publicName, + IODescriptor descriptor, + ContainerFlowInterfaceResolver.ExposedHandle exposedHandle, + boolean iterated) { + } + + public record ResolvedOutput( + String publicName, + IODescriptor descriptor, + ContainerFlowInterfaceResolver.ExposedHandle exposedHandle) { + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/types/IteratorContainerType.java b/src/main/java/it/cnr/isti/workflow/manager/containers/types/IteratorContainerType.java new file mode 100644 index 0000000..5941ea5 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/types/IteratorContainerType.java @@ -0,0 +1,37 @@ +package it.cnr.isti.workflow.manager.containers.types; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; + +@Component(IteratorContainerType.TYPE) +public class IteratorContainerType implements ContainerType { + + public static final String TYPE = "IteratorContainer"; + + @Override + public String getName() { + return TYPE; + } + + @Override + public String getDescription() { + return "A container node that iterates a subflow over one selected open input. The iterated container input is exposed as multiple and outputs are collected as lists."; + } + + @Override + public boolean validate() { + return true; + } + + @Override + public boolean isUserInteractive() { + return false; + } + + @Override + public Class> getContainerConfigurationClass() { + return IteratorContainerConfiguration.class; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java index efa4512..8e5a29c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ContainersController.java @@ -17,10 +17,10 @@ import org.springframework.web.bind.annotation.RestController; import org.springframework.web.server.ResponseStatusException; import it.cnr.isti.workflow.manager.containers.Container; -import it.cnr.isti.workflow.manager.containers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.ContainerSubFlowValidator; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; 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.types.ContainerType; import it.cnr.isti.workflow.manager.blocks.configurations.JsonSchemaProducer; import it.cnr.isti.workflow.manager.flows.model.FlowData; @@ -111,11 +111,11 @@ public class ContainersController { @SecurityRequirement(name = "bearerAuth") @Operation(summary = "Create container", description = "Creates a container from the provided container configuration.") public > Container create(@RequestBody C configuration) { - validateConfiguration(configuration); ContainerFactory factory = (ContainerFactory) containerFactories.stream() .filter(f -> f.getContainerType().equals(configuration.getContainerType())) .findFirst() .orElseThrow(() -> new IllegalArgumentException("Container factory not found for type: " + configuration.getContainerType())); + validateConfiguration(configuration, factory); return factory.create(configuration); } @@ -136,21 +136,25 @@ public class ContainersController { return new ContainerHandleDescriptor(handle.blockId(), handle.blockName(), handle.io()); } - private void validateConfiguration(ContainerConfiguration configuration) { + private > void validateConfiguration(C configuration, + ContainerFactory factory) { if (configuration == null) { return; } ContainerSubFlowValidator.ValidationResult validation = containerSubFlowValidator .validate(configuration.getSubFlow()); - if (validation.valid()) { - return; + if (!validation.valid()) { + String message = validation.errors().stream() + .map(error -> error.message()) + .collect(Collectors.joining(", ")); + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Invalid parameter subFlow: " + message); + } + try { + factory.create(configuration); + } catch (IllegalArgumentException exception) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, exception.getMessage(), exception); } - - String message = validation.errors().stream() - .map(error -> error.message()) - .collect(Collectors.joining(", ")); - throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Invalid parameter subFlow: " + message); } private ContainerConfigurationDescriptor toDescriptor(ContainerType containerType) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java index e5c50e8..9bab894 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -15,7 +15,7 @@ import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockCon import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; 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.ContainerConfiguration; import it.cnr.isti.workflow.manager.flows.model.Flow; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; @@ -130,7 +130,7 @@ public class ExecutionsService { collectHttpRequirement(requirements, block); } for (Container container : flow.getContainers() == null ? List.>of() : flow.getContainers()) { - if (container != null && container.getSpecificConfiguration() instanceof GenericContainerConfiguration containerConfiguration) { + if (container != null && container.getSpecificConfiguration() instanceof ContainerConfiguration containerConfiguration) { collectRequirements(requirements, containerConfiguration.getSubFlow()); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java index 5f9891f..a0e778c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java @@ -1,6 +1,5 @@ package it.cnr.isti.workflow.manager.executions.executors.containers; -import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.stream.Collectors; @@ -9,8 +8,8 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import it.cnr.isti.workflow.manager.containers.Container; -import it.cnr.isti.workflow.manager.containers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; import it.cnr.isti.workflow.manager.executions.ExecutionObject; import it.cnr.isti.workflow.manager.executions.ExecutionStatus; @@ -35,11 +34,7 @@ public class GenericContainerExecutor implements ContainerExecutor requirement.key()).distinct().toList()) { - if (authorizations.containsKey(key)) { - executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key)); - } - } + propagateAuthorizations(innerExecution, authorizations); Map inputPortsByName = ContainerFlowInterfaceResolver .getExposedInputs(configuration.getSubFlow()).stream() @@ -74,7 +69,7 @@ public class GenericContainerExecutor implements ContainerExecutor outputs = new LinkedHashMap<>(); + Map outputs = new java.util.LinkedHashMap<>(); for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : ContainerFlowInterfaceResolver .getExposedOutputs(configuration.getSubFlow())) { ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle(); @@ -86,6 +81,14 @@ public class GenericContainerExecutor implements ContainerExecutor authorizations) { + for (String key : innerExecution.getRequiredAuthorizations().stream().map(requirement -> requirement.key()).distinct().toList()) { + if (authorizations.containsKey(key)) { + executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key)); + } + } + } + @Override public Class getContainerType() { return GenericContainerType.class; 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 new file mode 100644 index 0000000..78e560c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -0,0 +1,124 @@ +package it.cnr.isti.workflow.manager.executions.executors.containers; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; +import it.cnr.isti.workflow.manager.executions.ExecutionObject; +import it.cnr.isti.workflow.manager.executions.ExecutionStatus; +import it.cnr.isti.workflow.manager.executions.ExecutionsService; +import it.cnr.isti.workflow.manager.executions.FieldKey; +import it.cnr.isti.workflow.manager.executions.steps.Input; + +@Component +public class IteratorContainerExecutor implements ContainerExecutor { + + private final ExecutionsService executionsService; + + public IteratorContainerExecutor(ExecutionsService executionsService) { + this.executionsService = executionsService; + } + + @Override + public Map execute(Container container, List inputs, + Map authorizations) { + IteratorContainerConfiguration configuration = (IteratorContainerConfiguration) container.getSpecificConfiguration(); + IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration); + IteratorContainerInterfaceResolver.ResolvedInput iteratedInput = resolution.inputsByPublicName() + .get(resolution.iteratedPublicName()); + Input runtimeIteratedInput = inputs.stream() + .filter(input -> input.getDescriptor().getName().equals(resolution.iteratedPublicName())) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException( + "Unknown IteratorContainer input port: " + resolution.iteratedPublicName())); + if (!(runtimeIteratedInput.getValue() instanceof List iterationValues)) { + throw new IllegalArgumentException( + "IteratorContainer input " + resolution.iteratedPublicName() + " expects a list value"); + } + + Map runtimeInputsByName = inputs.stream() + .collect(LinkedHashMap::new, (map, input) -> map.put(input.getDescriptor().getName(), input), Map::putAll); + + Map> collectedOutputs = new LinkedHashMap<>(); + for (IteratorContainerInterfaceResolver.ResolvedOutput output : resolution.resolvedOutputs()) { + collectedOutputs.put(output.publicName(), new ArrayList<>()); + } + + for (Object iterationValue : iterationValues) { + ExecutionObject innerExecution = executionsService.createExecution(container.getName() + " iteration", + configuration.getSubFlow()); + + propagateAuthorizations(innerExecution, authorizations); + + for (IteratorContainerInterfaceResolver.ResolvedInput resolvedInput : resolution.resolvedInputs()) { + Object value = resolvedInput.iterated() + ? iterationValue + : runtimeInputsByName.get(resolvedInput.publicName()).getValue(); + executionsService.prepareInput( + innerExecution.getId(), + resolvedInput.exposedHandle().handle().blockId(), + resolvedInput.exposedHandle().handle().io().getName(), + value); + } + + innerExecution = startAndWait(innerExecution); + + for (IteratorContainerInterfaceResolver.ResolvedOutput resolvedOutput : resolution.resolvedOutputs()) { + Object value = innerExecution.getContext().getResult() + .get(new FieldKey( + resolvedOutput.exposedHandle().handle().blockId(), + resolvedOutput.exposedHandle().handle().io().getName())); + collectedOutputs.get(resolvedOutput.publicName()).add(value); + } + } + + return new LinkedHashMap<>(collectedOutputs); + } + + private void propagateAuthorizations(ExecutionObject innerExecution, Map authorizations) { + for (String key : innerExecution.getRequiredAuthorizations().stream().map(requirement -> requirement.key()).distinct() + .toList()) { + if (authorizations.containsKey(key)) { + executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key)); + } + } + } + + private ExecutionObject startAndWait(ExecutionObject innerExecution) { + innerExecution = executionsService.startExecution(innerExecution.getId()); + while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) { + try { + Thread.sleep(25L); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Interrupted while executing IteratorContainer subflow", e); + } + innerExecution = executionsService.getExecution(innerExecution.getId()); + } + + if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) { + throw new IllegalStateException( + "IteratorContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet"); + } + if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) { + throw new IllegalStateException("IteratorContainer subflow failed: " + innerExecution.getContext().getErrors()); + } + if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) { + throw new IllegalStateException( + "IteratorContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus()); + } + return innerExecution; + } + + @Override + public Class getContainerType() { + return IteratorContainerType.class; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java index ffe5ede..487b49f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java @@ -17,9 +17,7 @@ import it.cnr.isti.workflow.manager.blocks.factories.ConditionalBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.ConditionalBlockType; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; -import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; import it.cnr.isti.workflow.manager.containers.factories.ContainerFactory; -import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.model.FlowNode; @@ -110,31 +108,52 @@ public class FlowDataValidator implements ConstraintValidator container) { - if (!(container.getSpecificConfiguration() instanceof GenericContainerConfiguration containerConfiguration)) { - return; + if (container.getSpecificConfiguration() == null) { + throw validationError(error("container", container.getId(), "specificConfiguration", + "Container is missing specificConfiguration")); } + + ContainerConfiguration containerConfiguration = container.getSpecificConfiguration(); FlowData subFlow = containerConfiguration.getSubFlow(); if (subFlow == null || subFlow.getNodes().isEmpty()) { throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "GenericContainer subFlow must contain at least one node")); + "Container subFlow must contain at least one node")); } List nestedNodes = subFlow.getNodes(); - if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container nestedContainer - && GenericContainerType.TYPE.equals(nestedContainer.getType().getName()))) { + if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container)) { throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Nested GenericContainer blocks are not supported")); + "Nested containers are not supported")); } if (nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) { throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Interactive blocks inside GenericContainer are not supported yet")); + "Interactive blocks inside containers are not supported yet")); } try { validateFlowData(subFlow); } catch (FlowValidationException e) { throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Invalid GenericContainer subFlow: " + e.getMessage())); + "Invalid container subFlow: " + e.getMessage())); + } + + final Container canonicalContainer; + try { + canonicalContainer = recreateContainer(containerConfiguration); + } catch (IllegalArgumentException exception) { + throw validationError(error("container", container.getId(), "specificConfiguration", exception.getMessage())); + } + if (!Objects.equals(container.getType().getName(), canonicalContainer.getType().getName())) { + throw validationError(error("container", container.getId(), "type", "Type does not match its configuration")); + } + if (!Objects.equals(container.getName(), canonicalContainer.getName())) { + throw validationError(error("container", container.getId(), "name", "Name does not match its configuration")); + } + if (!Objects.equals(container.getInputs(), canonicalContainer.getInputs())) { + throw validationError(error("container", container.getId(), "inputs", "Inputs do not match its configuration")); + } + if (!Objects.equals(container.getOutputs(), canonicalContainer.getOutputs())) { + throw validationError(error("container", container.getId(), "outputs", "Outputs do not match its configuration")); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java index cb834cb..f7e5738 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java @@ -11,8 +11,7 @@ import org.springframework.web.server.ResponseStatusException; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.containers.Container; -import it.cnr.isti.workflow.manager.containers.ContainerFlowInterfaceResolver; -import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.model.FlowNode; import jakarta.validation.ConstraintViolation; @@ -39,9 +38,11 @@ public class FlowExecutionValidator { public List collectErrors(FlowData flowData) { List errors = new ArrayList<>(); - errors.addAll(validator.validate(flowData).stream() - .flatMap(violation -> ValidationErrorCodec.decode(violation.getMessage()).stream()) - .toList()); + if (flowData != null) { + errors.addAll(validator.validate(flowData).stream() + .flatMap(violation -> ValidationErrorCodec.decode(violation.getMessage()).stream()) + .toList()); + } List nodes = flowData == null ? List.of() : flowData.getNodes(); for (FlowNode node : nodes) { @@ -64,7 +65,7 @@ public class FlowExecutionValidator { violation.getMessage())); } - if (container.getSpecificConfiguration() instanceof GenericContainerConfiguration containerConfiguration + if (container.getSpecificConfiguration() instanceof ContainerConfiguration containerConfiguration && containerConfiguration.getSubFlow() != null) { errors.addAll(collectErrors(containerConfiguration.getSubFlow()).stream() .map(error -> new ValidationError( diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java index 724f2a1..812dbd5 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/ContainersControllerTest.java @@ -23,7 +23,9 @@ import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; 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; import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; @@ -43,6 +45,7 @@ public class ContainersControllerTest { assertNotNull(types); assertFalse(types.isEmpty()); assertTrue(types.stream().anyMatch(type -> GenericContainerType.TYPE.equals(type.type()))); + assertTrue(types.stream().anyMatch(type -> IteratorContainerType.TYPE.equals(type.type()))); assertTrue(types.stream().allMatch(ContainersController.ContainerConfigurationDescriptor::hasExampleContainer)); assertTrue(types.stream().allMatch(type -> type.exampleContainerEndpoint().equals("/containers/types/" + type.type() + "/example"))); } @@ -62,6 +65,21 @@ public class ContainersControllerTest { assertTrue(container.getOutputs().isEmpty()); } + @Test + public void getIteratorContainerExampleForType() { + Container container = containersController.getExampleForType(IteratorContainerType.TYPE); + + assertNotNull(container); + assertEquals(IteratorContainerType.TYPE, container.getType().getName()); + assertEquals(IteratorContainerType.TYPE, container.getName()); + assertNotNull(container.getSpecificConfiguration()); + assertNotNull(container.getPosition()); + assertEquals(0, container.getPosition().x()); + assertEquals(0, container.getPosition().y()); + assertTrue(container.getInputs().isEmpty()); + assertTrue(container.getOutputs().isEmpty()); + } + @Test public void genericContainerSchemaContainsSecureSubFlowRetriever() { ContainersController.ContainerConfigurationDescriptor descriptor = containersController.getTypes().stream() @@ -219,4 +237,48 @@ public class ContainersControllerTest { assertTrue(exception.getReason().contains("Invalid parameter subFlow")); assertTrue(exception.getReason().contains("must contain at least one node")); } + + @Test + public void createIteratorContainerPluralizesIterationInputAndCollectsMultipleOutputs() { + Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("Analyze") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .prompt("Analyze ${{candidate}}") + .build()); + + Container container = containersController.create(IteratorContainerConfiguration.builder() + .name("Iterator") + .subFlow(FlowData.builder().block(internalBlock).build()) + .iterationInput("candidate") + .build()); + + assertNotNull(container); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidates") && input.isMultiple())); + assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response") && output.isMultiple())); + } + + @Test + public void createIteratorContainerRejectsUnknownIterationInput() { + Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("Analyze") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .prompt("Analyze ${{candidate}}") + .build()); + + ResponseStatusException exception = assertThrows(ResponseStatusException.class, + () -> containersController.create(IteratorContainerConfiguration.builder() + .name("Iterator") + .subFlow(FlowData.builder().block(internalBlock).build()) + .iterationInput("unknown") + .build())); + + assertEquals(400, exception.getStatusCode().value()); + assertTrue(exception.getReason().contains("iterationInput")); + } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java index 7e417d1..a6c283c 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java @@ -23,11 +23,15 @@ import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration; import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; import it.cnr.isti.workflow.manager.executions.steps.Step; import it.cnr.isti.workflow.manager.executions.steps.StepStatus; import it.cnr.isti.workflow.manager.containers.factories.GenericContainerFactory; +import it.cnr.isti.workflow.manager.containers.factories.IteratorContainerFactory; import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; import it.cnr.isti.workflow.manager.flows.ImportedFlow; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.FlowTestCreator; @@ -91,6 +95,9 @@ public class ExecutionTest { @Autowired GenericContainerFactory genericContainerFactory; + @Autowired + IteratorContainerFactory iteratorContainerFactory; + LLMDescriptor llmBrick = LLMDescriptor.builder() .provider("testProvider") .model("testModel") @@ -214,7 +221,7 @@ public class ExecutionTest { } } for (Container container : flowData.getContainers() == null ? List.>of() : flowData.getContainers()) { - if (container.getSpecificConfiguration() instanceof GenericContainerConfiguration containerConfiguration) { + if (container.getSpecificConfiguration() instanceof ContainerConfiguration containerConfiguration) { normalizeSeedFlowProviders(containerConfiguration.getSubFlow()); } } @@ -329,6 +336,41 @@ public class ExecutionTest { .anyMatch(key -> key.toString().equals(container.getId() + ":" + exposedOutputName))); } + @Test + public void iteratorContainerExecutesSubflowForEachIteratedValue() { + Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Internal LLM") + .llmDescriptor(llmBrick) + .prompt("Hello, ${{name}}!") + .build()); + + Container container = iteratorContainerFactory.create(IteratorContainerConfiguration.builder() + .name("Iterator") + .subFlow(FlowData.builder().block(internalBlock).build()) + .iterationInput("name") + .build()); + + FlowData flow = FlowData.builder().container(container).build(); + ExecutionObject execObject = executionsService.createExecution("Iterator flow", flow); + + executionsService.prepareInput(execObject.getId(), container.getId(), "names", List.of("Alice", "Bob")); + execObject = executionsService.startExecution(execObject.getId()); + while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { + try { + Thread.sleep(50); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + execObject = executionsService.getExecution(execObject.getId()); + } + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); + assertTrue(output instanceof List); + assertEquals(List.of("Hello, testModel!", "Hello, testModel!"), output); + } + @Test public void singleTextInputRejectsMultipleValues() { Flow flow = flowTestCreator.createFlowwithLLMUnpromptedWithConnection(llmBrick);