IteratorContainer added e small packaging refactoring

This commit is contained in:
Lucio Lelii 2026-03-17 10:15:45 +01:00
parent bdbe5654df
commit 56424d2f0f
16 changed files with 622 additions and 41 deletions

View File

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

View File

@ -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<IteratorContainerType> {
@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<IteratorContainerType> getContainerType() {
return IteratorContainerType.class;
}
public static IteratorContainerConfiguration empty() {
return new IteratorContainerConfiguration(
IteratorContainerType.TYPE,
FlowData.builder().build(),
null);
}
}

View File

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

View File

@ -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<IteratorContainerType, IteratorContainerConfiguration> {
private static final Position DEFAULT_POSITION = new Position(0, 0);
private final IteratorContainerType containerType;
public IteratorContainerFactory(IteratorContainerType containerType) {
this.containerType = containerType;
}
@Override
public Container<IteratorContainerType> create(IteratorContainerConfiguration configuration) {
ContainerFlowInterface exposedInterface = isEmpty(configuration)
? new ContainerFlowInterface(java.util.List.of(), java.util.List.of())
: IteratorContainerInterfaceResolver.resolve(configuration);
return Container.<IteratorContainerType>builder()
.inputs(exposedInterface.inputs())
.outputs(exposedInterface.outputs())
.specificConfiguration(configuration)
.type(containerType)
.position(DEFAULT_POSITION)
.build();
}
@Override
public Container<IteratorContainerType> createEmpty() {
return create(IteratorContainerConfiguration.empty());
}
@Override
public Class<IteratorContainerType> getContainerType() {
return IteratorContainerType.class;
}
private boolean isEmpty(IteratorContainerConfiguration configuration) {
return configuration == null
|| configuration.getSubFlow() == null
|| configuration.getSubFlow().getNodes().isEmpty();
}
}

View File

@ -1,4 +1,4 @@
package it.cnr.isti.workflow.manager.containers;
package it.cnr.isti.workflow.manager.containers.iresolvers;
import java.util.List;

View File

@ -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<ExposedHandle> openInputs = getExposedInputs(subFlow);

View File

@ -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<ContainerFlowInterfaceResolver.ExposedHandle> exposedInputs = ContainerFlowInterfaceResolver
.getExposedInputs(configuration.getSubFlow());
List<ContainerFlowInterfaceResolver.ExposedHandle> 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<ResolvedInput> 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<ResolvedOutput> 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<String> reservedNames) {
Set<String> 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<ResolvedInput> resolvedInputs,
List<ResolvedOutput> resolvedOutputs,
ContainerFlowInterfaceResolver.ExposedHandle iteratedHandle,
String iteratedPublicName) {
public List<IODescriptor> inputs() {
return resolvedInputs.stream().map(ResolvedInput::descriptor).toList();
}
public List<IODescriptor> outputs() {
return resolvedOutputs.stream().map(ResolvedOutput::descriptor).toList();
}
public Map<String, ResolvedInput> 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) {
}
}

View File

@ -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<? extends ContainerConfiguration<?>> getContainerConfigurationClass() {
return IteratorContainerConfiguration.class;
}
}

View File

@ -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 <T extends ContainerType, C extends ContainerConfiguration<T>> Container<T> create(@RequestBody C configuration) {
validateConfiguration(configuration);
ContainerFactory<T, C> factory = (ContainerFactory<T, C>) 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 <T extends ContainerType, C extends ContainerConfiguration<T>> void validateConfiguration(C configuration,
ContainerFactory<T, C> 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) {

View File

@ -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.<Container<?>>of() : flow.getContainers()) {
if (container != null && container.getSpecificConfiguration() instanceof GenericContainerConfiguration containerConfiguration) {
if (container != null && container.getSpecificConfiguration() instanceof ContainerConfiguration<?> containerConfiguration) {
collectRequirements(requirements, containerConfiguration.getSubFlow());
}
}

View File

@ -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<GenericContai
ExecutionObject innerExecution = executionsService.createExecution(container.getName() + " subflow", configuration.getSubFlow());
for (String key : innerExecution.getRequiredAuthorizations().stream().map(requirement -> requirement.key()).distinct().toList()) {
if (authorizations.containsKey(key)) {
executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key));
}
}
propagateAuthorizations(innerExecution, authorizations);
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
.getExposedInputs(configuration.getSubFlow()).stream()
@ -74,7 +69,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
throw new IllegalStateException("GenericContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus());
}
Map<String, Object> outputs = new LinkedHashMap<>();
Map<String, Object> 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<GenericContai
return outputs;
}
private void propagateAuthorizations(ExecutionObject innerExecution, Map<String, Object> 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<GenericContainerType> getContainerType() {
return GenericContainerType.class;

View File

@ -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<IteratorContainerType> {
private final ExecutionsService executionsService;
public IteratorContainerExecutor(ExecutionsService executionsService) {
this.executionsService = executionsService;
}
@Override
public Map<String, Object> execute(Container<IteratorContainerType> container, List<Input> inputs,
Map<String, Object> 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<String, Input> runtimeInputsByName = inputs.stream()
.collect(LinkedHashMap::new, (map, input) -> map.put(input.getDescriptor().getName(), input), Map::putAll);
Map<String, List<Object>> 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<String, Object> 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<IteratorContainerType> getContainerType() {
return IteratorContainerType.class;
}
}

View File

@ -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<ValidFlowStructure
}
private void validateContainer(Container<?> 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<FlowNode> 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"));
}
}

View File

@ -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<ValidationError> collectErrors(FlowData flowData) {
List<ValidationError> 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<FlowNode> 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(

View File

@ -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<IteratorContainerType> 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<LLMBlockType> internalBlock = blocksController.create(LLMBlockConfiguration.builder()
.name("Analyze")
.llmDescriptor(LLMDescriptor.builder()
.provider("testProvider")
.model("testModel")
.build())
.prompt("Analyze ${{candidate}}")
.build());
Container<IteratorContainerType> 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<LLMBlockType> 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"));
}
}

View File

@ -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.<Container<?>>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<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
.name("Internal LLM")
.llmDescriptor(llmBrick)
.prompt("Hello, ${{name}}!")
.build());
Container<IteratorContainerType> 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);