Refine container interfaces and conditional schema metadata

This commit is contained in:
Lucio Lelii 2026-03-16 15:42:05 +01:00
parent a44d11fb1d
commit b268014d8c
18 changed files with 314 additions and 381 deletions

View File

@ -6,6 +6,7 @@ import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.Structural;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDependency;
import it.cnr.isti.workflow.manager.configurations.annotations.UiRequiredWhen;
import it.cnr.isti.workflow.manager.blocks.types.ConditionalBlockType;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import jakarta.validation.Valid;
@ -35,9 +36,11 @@ public class ConditionalBlockConfiguration extends BlockConfiguration<Conditiona
@Valid
@UiDependency(field = "useLlm", equals = "true", group = "llm")
@UiRequiredWhen(field = "useLlm", equals = "true")
LLMDescriptor llmDescriptor;
@UiDependency(field = "useLlm", equals = "true", group = "llm")
@UiRequiredWhen(field = "useLlm", equals = "true")
@Structural
@LongText(
placeholder = "Add the prompt the LLM should use to decide true or false",

View File

@ -28,6 +28,7 @@ import it.cnr.isti.workflow.manager.configurations.annotations.Structural;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDependency;
import it.cnr.isti.workflow.manager.configurations.annotations.DynamicSchema;
import it.cnr.isti.workflow.manager.configurations.annotations.ConfigurableAsInput;
import it.cnr.isti.workflow.manager.configurations.annotations.UiRequiredWhen;
@Component
public class JsonSchemaProducer {
@ -49,12 +50,14 @@ public class JsonSchemaProducer {
Map<Class<?>, Map<String, LongText>> longTextMap = collectLongTextMetadata(type);
Map<Class<?>, Map<String, Structural>> structuralMap = collectStructuralMetadata(type);
Map<Class<?>, Map<String, UiDependency>> uiDependencyMap = collectUiDependencyMetadata(type);
Map<Class<?>, Map<String, UiRequiredWhen>> uiRequiredWhenMap = collectUiRequiredWhenMetadata(type);
Map<Class<?>, Map<String, ConfigurableAsInput>> configurableAsInputMap = collectConfigurableAsInputMetadata(type);
applyRetrieverMetadata(root, type, retrieverMap.getOrDefault(type, Map.of()));
applyDynamicSchemaMetadata(root, dynamicSchemaMap.getOrDefault(type, Map.of()));
applyLongTextMetadata(root, longTextMap.getOrDefault(type, Map.of()));
applyStructuralMetadata(root, structuralMap.getOrDefault(type, Map.of()));
applyUiDependencyMetadata(root, uiDependencyMap.getOrDefault(type, Map.of()));
applyUiRequiredWhenMetadata(root, uiRequiredWhenMap.getOrDefault(type, Map.of()));
applyConfigurableAsInputMetadata(root, configurableAsInputMap.getOrDefault(type, Map.of()));
JsonNode definitionsNode = root.has("definitions") ? root.get("definitions") : root.get("$defs");
@ -64,6 +67,7 @@ public class JsonSchemaProducer {
metadataClasses.addAll(longTextMap.keySet());
metadataClasses.addAll(structuralMap.keySet());
metadataClasses.addAll(uiDependencyMap.keySet());
metadataClasses.addAll(uiRequiredWhenMap.keySet());
metadataClasses.addAll(configurableAsInputMap.keySet());
for (Entry<String, JsonNode> entry : iterable(definitions.fields())) {
if (!(entry.getValue() instanceof ObjectNode classSchema)) {
@ -76,6 +80,7 @@ public class JsonSchemaProducer {
applyLongTextMetadata(classSchema, longTextMap.get(matchedClass));
applyStructuralMetadata(classSchema, structuralMap.get(matchedClass));
applyUiDependencyMetadata(classSchema, uiDependencyMap.get(matchedClass));
applyUiRequiredWhenMetadata(classSchema, uiRequiredWhenMap.get(matchedClass));
applyConfigurableAsInputMetadata(classSchema, configurableAsInputMap.get(matchedClass));
}
}
@ -432,6 +437,75 @@ public class JsonSchemaProducer {
}
}
private Map<Class<?>, Map<String, UiRequiredWhen>> collectUiRequiredWhenMetadata(Class<?> rootClass) {
Map<Class<?>, Map<String, UiRequiredWhen>> result = new HashMap<>();
Set<Class<?>> visited = new HashSet<>();
Queue<Class<?>> queue = new ArrayDeque<>();
queue.add(rootClass);
while (!queue.isEmpty()) {
Class<?> current = queue.poll();
if (current == null || !visited.add(current) || isTerminalType(current)) {
continue;
}
Map<String, UiRequiredWhen> metadata = new LinkedHashMap<>();
for (Field field : current.getDeclaredFields()) {
UiRequiredWhen annotation = field.getAnnotation(UiRequiredWhen.class);
if (annotation != null) {
metadata.put(field.getName(), annotation);
}
enqueueRelatedTypes(queue, field.getGenericType(), field.getType());
}
if (current.isRecord()) {
for (RecordComponent component : current.getRecordComponents()) {
UiRequiredWhen annotation = component.getAnnotation(UiRequiredWhen.class);
if (annotation != null) {
metadata.put(component.getName(), annotation);
}
enqueueRelatedTypes(queue, component.getGenericType(), component.getType());
}
}
if (!metadata.isEmpty()) {
result.put(current, metadata);
}
}
return result;
}
private void applyUiRequiredWhenMetadata(ObjectNode classSchema, Map<String, UiRequiredWhen> metadata) {
if (metadata == null || metadata.isEmpty()) {
return;
}
JsonNode propsNode = classSchema.get("properties");
if (!(propsNode instanceof ObjectNode properties)) {
return;
}
for (Entry<String, UiRequiredWhen> entry : metadata.entrySet()) {
JsonNode propNode = properties.get(entry.getKey());
if (!(propNode instanceof ObjectNode propertySchema)) {
continue;
}
UiRequiredWhen dependency = entry.getValue();
ObjectNode requiredWhen = propertySchema.putObject("x-ui-required-when");
requiredWhen.put("field", dependency.field());
if (!dependency.equals().isBlank()) {
requiredWhen.put("equals", dependency.equals());
}
if (dependency.equalsAny().length > 0) {
ArrayNode equalsAny = requiredWhen.putArray("in");
for (String value : dependency.equalsAny()) {
equalsAny.add(value);
}
}
}
}
private Map<Class<?>, Map<String, ConfigurableAsInput>> collectConfigurableAsInputMetadata(Class<?> rootClass) {
Map<Class<?>, Map<String, ConfigurableAsInput>> result = new HashMap<>();
Set<Class<?>> visited = new HashSet<>();

View File

@ -1,15 +0,0 @@
package it.cnr.isti.workflow.manager.blocks.configurations.retrievers;
import java.util.List;
import java.util.Map;
public interface DynamicFieldRetriever {
String getCategory();
List<String> retrieve(String parameter, Map<String, String> params);
default boolean isRequired(String parameter, Map<String, String> params) {
return false;
}
}

View File

@ -1,76 +0,0 @@
package it.cnr.isti.workflow.manager.blocks.configurations.retrievers;
import java.util.List;
import java.util.Map;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
@Component
public class LLMFieldRetriever implements DynamicFieldRetriever {
@Autowired
private Map<String, LLMProvider> llmProviders;
@Override
public String getCategory() {
return "LLM";
}
@Override
public List<String> retrieve(String parameter, Map<String, String> params) {
return switch (parameter) {
case "providers" -> llmProviders.values().stream()
.map(LLMProvider::getName)
.distinct()
.sorted()
.toList();
case "models" -> {
LLMProvider provider = resolveProvider(params.get("provider"));
yield provider.getRegisteredModels();
}
case "authorization" -> List.of();
default -> throw new ResponseStatusException(HttpStatus.NOT_FOUND,
"Unknown LLM retriever parameter: " + parameter);
};
}
@Override
public boolean isRequired(String parameter, Map<String, String> params) {
return switch (parameter) {
case "providers" -> true;
case "models" -> {
String provider = params.get("provider");
yield provider != null && !provider.isBlank();
}
case "authorization" -> resolveProviderNullable(params.get("provider"))
.map(LLMProvider::requiresAuthorization)
.orElse(false);
default -> false;
};
}
private LLMProvider resolveProvider(String providerName) {
if (providerName == null || providerName.isBlank()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Missing required parameter: provider");
}
return resolveProviderNullable(providerName)
.orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND,
"Provider not found: " + providerName));
}
private java.util.Optional<LLMProvider> resolveProviderNullable(String providerName) {
if (providerName == null || providerName.isBlank()) {
return java.util.Optional.empty();
}
LLMProvider llmProvider = llmProviders.get(providerName);
if (llmProvider != null) {
return java.util.Optional.of(llmProvider);
}
return llmProviders.values().stream().filter(p -> p.getName().equals(providerName)).findFirst();
}
}

View File

@ -1,37 +0,0 @@
package it.cnr.isti.workflow.manager.blocks.configurations.retrievers;
import java.util.List;
import java.util.Map;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.mcp.MCPServersProvider;
@Component
public class MCPServersFieldRetriever implements DynamicFieldRetriever {
private final MCPServersProvider mcpServersProvider;
public MCPServersFieldRetriever(MCPServersProvider mcpServersProvider) {
this.mcpServersProvider = mcpServersProvider;
}
@Override
public String getCategory() {
return "MCPServers";
}
@Override
public List<String> retrieve(String parameter, Map<String, String> params) {
return switch (parameter) {
case "servers" -> mcpServersProvider.getServers().stream()
.map(MCPServersProvider.MCPServerDefinition::id)
.sorted()
.toList();
default -> throw new ResponseStatusException(HttpStatus.NOT_FOUND,
"Unknown MCPServers retriever parameter: " + parameter);
};
}
}

View File

@ -0,0 +1,16 @@
package it.cnr.isti.workflow.manager.configurations.annotations;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
@Target({ ElementType.FIELD, ElementType.RECORD_COMPONENT })
@Retention(RetentionPolicy.RUNTIME)
public @interface UiRequiredWhen {
String field();
String equals() default "";
String[] equalsAny() default {};
}

View File

@ -4,8 +4,9 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
import java.util.stream.Collectors;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
@ -19,31 +20,26 @@ public final class ContainerFlowInterfaceResolver {
public static ContainerFlowInterface resolve(GenericContainerConfiguration configuration) {
FlowData subFlow = configuration == null ? null : configuration.getSubFlow();
List<OpenHandle> openInputs = getOpenInputs(subFlow);
List<OpenHandle> openOutputs = getOpenOutputs(subFlow);
Map<String, OpenHandle> openInputsByKey = indexOpenHandles(openInputs);
Map<String, OpenHandle> openOutputsByKey = indexOpenHandles(openOutputs);
List<GenericContainerConfiguration.InputPort> publicInputs = configuration == null
|| configuration.getPublicInputs() == null ? List.of() : configuration.getPublicInputs();
List<GenericContainerConfiguration.OutputPort> publicOutputs = configuration == null
|| configuration.getPublicOutputs() == null ? List.of() : configuration.getPublicOutputs();
List<IODescriptor> inputs = publicInputs.isEmpty()
? openInputs.stream().map(handle -> cloneDescriptor(handle.io().getName(), handle.io())).toList()
: publicInputs.stream()
.map(port -> toPublicInput(port, openInputsByKey))
.filter(Objects::nonNull)
.toList();
List<IODescriptor> outputs = publicOutputs.isEmpty()
? openOutputs.stream().map(handle -> cloneDescriptor(handle.io().getName(), handle.io())).toList()
: publicOutputs.stream()
.map(port -> toPublicOutput(port, openOutputsByKey))
.filter(Objects::nonNull)
.toList();
List<ExposedHandle> openInputs = getExposedInputs(subFlow);
List<ExposedHandle> openOutputs = getExposedOutputs(subFlow);
List<IODescriptor> inputs = openInputs.stream()
.map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io()))
.toList();
List<IODescriptor> outputs = openOutputs.stream()
.map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io()))
.toList();
return new ContainerFlowInterface(inputs, outputs);
}
public static List<ExposedHandle> getExposedInputs(FlowData subFlow) {
return expose(getOpenInputs(subFlow));
}
public static List<ExposedHandle> getExposedOutputs(FlowData subFlow) {
return expose(getOpenOutputs(subFlow));
}
public static List<OpenHandle> getOpenInputs(FlowData subFlow) {
List<FlowNode> blocks = subFlow == null ? List.of() : subFlow.getNodes();
List<Connection> connections = subFlow == null || subFlow.getConnections() == null ? List.of() : subFlow.getConnections();
@ -86,34 +82,6 @@ public final class ContainerFlowInterfaceResolver {
.findFirst();
}
private static Map<String, OpenHandle> indexOpenHandles(List<OpenHandle> handles) {
Map<String, OpenHandle> indexed = new LinkedHashMap<>();
for (OpenHandle handle : handles) {
indexed.put(key(handle.blockId(), handle.io().getName()), handle);
}
return indexed;
}
private static IODescriptor toPublicInput(GenericContainerConfiguration.InputPort port, Map<String, OpenHandle> openInputsByKey) {
OpenHandle handle = openInputsByKey.get(key(port.getTargetBlockId(), port.getTargetInputName()));
if (handle == null) {
return null;
}
return cloneDescriptor(port.getName(), handle.io());
}
private static IODescriptor toPublicOutput(GenericContainerConfiguration.OutputPort port, Map<String, OpenHandle> openOutputsByKey) {
OpenHandle handle = openOutputsByKey.get(key(port.getSourceBlockId(), port.getSourceOutputName()));
if (handle == null) {
return null;
}
return cloneDescriptor(port.getName(), handle.io());
}
private static String key(String blockId, String handleName) {
return blockId + "::" + handleName;
}
private static boolean isTargeted(List<Connection> connections, String blockId, String inputName) {
return connections.stream()
.anyMatch(connection -> blockId.equals(connection.getTargetId()) && inputName.equals(connection.getTargetName()));
@ -128,6 +96,44 @@ public final class ContainerFlowInterfaceResolver {
return new IODescriptor(publicName, source.getType(), source.isMultiple(), source.getValueKinds());
}
private static List<ExposedHandle> expose(List<OpenHandle> handles) {
Map<String, Long> names = handles.stream()
.collect(Collectors.groupingBy(handle -> handle.io().getName(), LinkedHashMap::new, Collectors.counting()));
List<OpenHandle> colliding = handles.stream()
.filter(handle -> names.getOrDefault(handle.io().getName(), 0L) > 1)
.toList();
Map<String, Long> namesWithNode = colliding.stream()
.collect(Collectors.groupingBy(
handle -> qualifyWithNodeName(handle),
LinkedHashMap::new,
Collectors.counting()));
List<ExposedHandle> exposed = new ArrayList<>();
for (OpenHandle handle : handles) {
String publicName = handle.io().getName();
if (names.getOrDefault(publicName, 0L) > 1) {
publicName = qualifyWithNodeName(handle);
if (namesWithNode.getOrDefault(publicName, 0L) > 1) {
publicName = qualifyWithNodeId(handle);
}
}
exposed.add(new ExposedHandle(publicName, handle));
}
return List.copyOf(exposed);
}
private static String qualifyWithNodeName(OpenHandle handle) {
return handle.blockName() + "." + handle.io().getName();
}
private static String qualifyWithNodeId(OpenHandle handle) {
return handle.blockId() + "." + handle.io().getName();
}
public record OpenHandle(String blockId, String blockName, IODescriptor io) {
}
public record ExposedHandle(String publicName, OpenHandle handle) {
}
}

View File

@ -6,7 +6,6 @@ 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.configurations.ContainerConfiguration;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import lombok.AllArgsConstructor;
import jakarta.validation.Valid;
import lombok.Builder;
import lombok.EqualsAndHashCode;
@ -30,23 +29,10 @@ public class GenericContainerConfiguration extends ContainerConfiguration<Generi
@JsonProperty(required = false)
private FlowData subFlow;
@Structural
@Valid
@JsonProperty(required = false)
private java.util.List<InputPort> publicInputs;
@Structural
@Valid
@JsonProperty(required = false)
private java.util.List<OutputPort> publicOutputs;
@Builder
public GenericContainerConfiguration(@NonNull String name, FlowData subFlow, java.util.List<InputPort> publicInputs,
java.util.List<OutputPort> publicOutputs) {
public GenericContainerConfiguration(@NonNull String name, FlowData subFlow) {
super(name);
this.subFlow = subFlow == null ? FlowData.builder().build() : subFlow;
this.publicInputs = publicInputs == null ? java.util.List.of() : java.util.List.copyOf(publicInputs);
this.publicOutputs = publicOutputs == null ? java.util.List.of() : java.util.List.copyOf(publicOutputs);
}
@Override
@ -57,34 +43,6 @@ public class GenericContainerConfiguration extends ContainerConfiguration<Generi
public static GenericContainerConfiguration empty() {
return new GenericContainerConfiguration(
GenericContainerType.TYPE,
FlowData.builder().build(),
java.util.List.of(),
java.util.List.of());
}
@Getter
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@AllArgsConstructor
@Builder
public static class InputPort {
@NonNull
private String name;
@NonNull
private String targetBlockId;
@NonNull
private String targetInputName;
}
@Getter
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@AllArgsConstructor
@Builder
public static class OutputPort {
@NonNull
private String name;
@NonNull
private String sourceBlockId;
@NonNull
private String sourceOutputName;
FlowData.builder().build());
}
}

View File

@ -3,7 +3,6 @@ package it.cnr.isti.workflow.manager.containers;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.springframework.beans.factory.annotation.Autowired;
@ -39,29 +38,16 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
}
}
List<GenericContainerConfiguration.InputPort> configuredInputPorts = configuration.getPublicInputs() == null
? List.of()
: configuration.getPublicInputs();
List<GenericContainerConfiguration.OutputPort> configuredOutputPorts = configuration.getPublicOutputs() == null
? List.of()
: configuration.getPublicOutputs();
Map<String, GenericContainerConfiguration.InputPort> inputPortsByName = configuredInputPorts.isEmpty()
? ContainerFlowInterfaceResolver.getOpenInputs(configuration.getSubFlow()).stream()
.collect(Collectors.toMap(
handle -> handle.io().getName(),
handle -> GenericContainerConfiguration.InputPort.builder()
.name(handle.io().getName())
.targetBlockId(handle.blockId())
.targetInputName(handle.io().getName())
.build()))
: configuredInputPorts.stream()
.collect(Collectors.toMap(GenericContainerConfiguration.InputPort::getName, Function.identity()));
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
.getExposedInputs(configuration.getSubFlow()).stream()
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed));
for (Input input : inputs) {
GenericContainerConfiguration.InputPort port = inputPortsByName.get(input.getDescriptor().getName());
if (port == null) {
ContainerFlowInterfaceResolver.ExposedHandle exposedHandle = inputPortsByName.get(input.getDescriptor().getName());
if (exposedHandle == null) {
throw new IllegalArgumentException("Unknown GenericContainer input port: " + input.getDescriptor().getName());
}
executionsService.prepareInput(innerExecution.getId(), port.getTargetBlockId(), port.getTargetInputName(), input.getValue());
ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle();
executionsService.prepareInput(innerExecution.getId(), handle.blockId(), handle.io().getName(), input.getValue());
}
innerExecution = executionsService.startExecution(innerExecution.getId());
@ -86,20 +72,12 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
}
Map<String, Object> outputs = new LinkedHashMap<>();
List<GenericContainerConfiguration.OutputPort> outputPorts = configuredOutputPorts.isEmpty()
? ContainerFlowInterfaceResolver.getOpenOutputs(configuration.getSubFlow()).stream()
.map(handle -> GenericContainerConfiguration.OutputPort.builder()
.name(handle.io().getName())
.sourceBlockId(handle.blockId())
.sourceOutputName(handle.io().getName())
.build())
.toList()
: configuredOutputPorts;
for (GenericContainerConfiguration.OutputPort port : outputPorts) {
Object value = innerExecution.getContext().getResult()
.get(new FieldKey(port.getSourceBlockId(), port.getSourceOutputName()));
for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : ContainerFlowInterfaceResolver
.getExposedOutputs(configuration.getSubFlow())) {
ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle();
Object value = innerExecution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName()));
if (value != null) {
outputs.put(port.getName(), value);
outputs.put(exposedHandle.publicName(), value);
}
}
return outputs;

View File

@ -3,11 +3,14 @@ package it.cnr.isti.workflow.manager.containers;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Position;
import it.cnr.isti.workflow.manager.containers.factories.ContainerFactory;
@Component
public class GenericContainerFactory implements ContainerFactory<GenericContainerType, GenericContainerConfiguration> {
private static final Position DEFAULT_POSITION = new Position(0, 0);
@Autowired
private GenericContainerType containerType;
@ -19,6 +22,7 @@ public class GenericContainerFactory implements ContainerFactory<GenericContaine
.outputs(exposedInterface.outputs())
.specificConfiguration(configuration)
.type(containerType)
.position(DEFAULT_POSITION)
.build();
}

View File

@ -40,6 +40,20 @@ public class FlowController {
public List<FlowView> getAllFlows(@AuthenticationPrincipal LoginEntity userDetails) {
return flowService.getAllFlows(userDetails.getUsername());
}
@GetMapping("/{id}")
@Operation(summary = "Get flow", description = "Returns a single flow visible to the authenticated user.")
public ResponseEntity<FlowView> getFlow(@PathVariable String id, @AuthenticationPrincipal LoginEntity userDetails) {
try {
return ResponseEntity.ok(flowService.getFlow(id, userDetails.getUsername()));
} catch (FlowNotFoundException e) {
logger.warn("Get flow {} failed for user {}: not found", id, userDetails.getUsername());
return ResponseEntity.status(HttpStatus.NOT_FOUND).build();
} catch (FlowAccessDeniedException e) {
logger.warn("Get flow {} forbidden for user {}", id, userDetails.getUsername());
return ResponseEntity.status(HttpStatus.FORBIDDEN).build();
}
}
@PostMapping

View File

@ -71,6 +71,18 @@ public class FlowService {
return flowRepository.findFlowsByOwnerOrPublic(owner).stream().map(this::toView).toList();
}
public FlowView getFlow(String id, String owner) {
FlowEntity entity = flowRepository.findById(id)
.orElseThrow(() -> new FlowNotFoundException(id));
if (!entity.getOwner().equals(owner) && !entity.isPublished()) {
logger.warn("User {} attempted to access flow {} owned by {}", owner, id, entity.getOwner());
throw new FlowAccessDeniedException();
}
return toView(entity);
}
private void validateFlow(FlowCreateRequest request) {
Set<ConstraintViolation<FlowCreateRequest>> violations = validator.validate(request);
if (violations.isEmpty()) {

View File

@ -1,7 +1,6 @@
package it.cnr.isti.workflow.manager.flows.validation;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
@ -13,7 +12,6 @@ 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.ContainerFlowInterfaceResolver.OpenHandle;
import it.cnr.isti.workflow.manager.containers.GenericContainerConfiguration;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
@ -75,86 +73,10 @@ public class FlowExecutionValidator {
"specificConfiguration.subFlow",
error.message()))
.toList());
errors.addAll(validateContainerPorts(container.getId(), containerConfiguration));
}
}
}
return errors;
}
private List<ValidationError> validateContainerPorts(String containerId, GenericContainerConfiguration configuration) {
List<ValidationError> errors = new ArrayList<>();
List<GenericContainerConfiguration.InputPort> publicInputs = configuration.getPublicInputs() == null
? List.of()
: configuration.getPublicInputs();
List<GenericContainerConfiguration.OutputPort> publicOutputs = configuration.getPublicOutputs() == null
? List.of()
: configuration.getPublicOutputs();
if (publicInputs.isEmpty()) {
errors.addAll(validateAutoExposedInputNames(containerId, configuration.getSubFlow()));
}
if (publicOutputs.isEmpty()) {
errors.addAll(validateAutoExposedOutputNames(containerId, configuration.getSubFlow()));
}
Set<String> publicInputNames = new HashSet<>();
for (GenericContainerConfiguration.InputPort inputPort : publicInputs) {
if (inputPort == null) {
continue;
}
if (!publicInputNames.add(inputPort.getName())) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.publicInputs",
"GenericContainer public input names must be unique"));
}
if (ContainerFlowInterfaceResolver.findOpenInput(configuration.getSubFlow(), inputPort.getTargetBlockId(),
inputPort.getTargetInputName()).isEmpty()) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.publicInputs",
"GenericContainer input port '" + inputPort.getName() + "' must target an open subflow input"));
}
}
Set<String> publicOutputNames = new HashSet<>();
for (GenericContainerConfiguration.OutputPort outputPort : publicOutputs) {
if (outputPort == null) {
continue;
}
if (!publicOutputNames.add(outputPort.getName())) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.publicOutputs",
"GenericContainer public output names must be unique"));
}
if (ContainerFlowInterfaceResolver.findOpenOutput(configuration.getSubFlow(), outputPort.getSourceBlockId(),
outputPort.getSourceOutputName()).isEmpty()) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.publicOutputs",
"GenericContainer output port '" + outputPort.getName() + "' must target an open subflow output"));
}
}
return errors;
}
private List<ValidationError> validateAutoExposedInputNames(String containerId, FlowData subFlow) {
List<ValidationError> errors = new ArrayList<>();
Set<String> names = new HashSet<>();
for (OpenHandle handle : ContainerFlowInterfaceResolver.getOpenInputs(subFlow)) {
String name = handle.io().getName();
if (!names.add(name)) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.subFlow",
"GenericContainer auto-exposed input names must be unique. Duplicate open input: " + name));
}
}
return errors;
}
private List<ValidationError> validateAutoExposedOutputNames(String containerId, FlowData subFlow) {
List<ValidationError> errors = new ArrayList<>();
Set<String> names = new HashSet<>();
for (OpenHandle handle : ContainerFlowInterfaceResolver.getOpenOutputs(subFlow)) {
String name = handle.io().getName();
if (!names.add(name)) {
errors.add(new ValidationError("container", containerId, "specificConfiguration.subFlow",
"GenericContainer auto-exposed output names must be unique. Duplicate open output: " + name));
}
}
return errors;
}
}

View File

@ -265,7 +265,7 @@
],
"outputs": [
{
"name": "analysis",
"name": "response",
"type": "TEXT",
"multiple": false,
"valueKinds": [
@ -279,20 +279,6 @@
"specificConfiguration": {
"type": "GenericContainerConfiguration",
"name": "candidate-analysis-container",
"publicInputs": [
{
"name": "candidate",
"targetBlockId": "7c88d1a3-84f6-43f1-80bc-f9bfe4d53003",
"targetInputName": "candidate"
}
],
"publicOutputs": [
{
"name": "analysis",
"sourceBlockId": "7c88d1a3-84f6-43f1-80bc-f9bfe4d53003",
"sourceOutputName": "response"
}
],
"subFlow": {
"blocks": [
{
@ -361,11 +347,11 @@
{
"id": "c31a0a91-d0d5-4215-9d5f-b4b0af2d4c02",
"sourceId": "3df53d5b-4f11-4387-843b-4c5a8a6a2002",
"sourceName": "analysis",
"sourceName": "response",
"targetId": "9e22f57b-b264-44c8-bf7c-64f8a7d63003",
"targetName": "analysis"
}
]
}
}
]
]

View File

@ -18,6 +18,8 @@ import com.fasterxml.jackson.databind.JsonNode;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.IOCapabilityType;
import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType;
import it.cnr.isti.workflow.manager.blocks.types.ConditionalBlockType;
import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.IterationBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
@ -273,6 +275,24 @@ public class BlocksControllerTest {
assertEquals("authorizationType", authHeader.path("x-ui-visible-when").path("field").asText());
}
@Test
public void conditionalBlockSchemaContainsConditionalRequiredMetadata() {
BlockConfigurationDescriptor descriptor = blocksController
.getConfigurationDescriptorForType(ConditionalBlockType.TYPE);
assertNotNull(descriptor.schema());
assertTrue(descriptor.schema() instanceof JsonNode);
JsonNode schema = (JsonNode) descriptor.schema();
JsonNode llmDescriptor = schema.path("properties").path("llmDescriptor");
JsonNode prompt = schema.path("properties").path("prompt");
assertEquals("useLlm", llmDescriptor.path("x-ui-required-when").path("field").asText());
assertEquals("true", llmDescriptor.path("x-ui-required-when").path("equals").asText());
assertEquals("useLlm", prompt.path("x-ui-required-when").path("field").asText());
assertEquals("true", prompt.path("x-ui-required-when").path("equals").asText());
}
@Test
public void getMCPBridgeExampleForType() {
Block<MCPBridgeBlockType> block = blocksController.getExampleForType(MCPBridgeBlockType.TYPE);

View File

@ -17,6 +17,7 @@ import org.springframework.web.server.ResponseStatusException;
import com.fasterxml.jackson.databind.JsonNode;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.Position;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType;
import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
@ -54,6 +55,9 @@ public class ContainersControllerTest {
assertEquals(GenericContainerType.TYPE, container.getType().getName());
assertEquals(GenericContainerType.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());
}
@ -79,7 +83,7 @@ public class ContainersControllerTest {
}
@Test
public void createGenericContainerExposesConfiguredPublicPorts() {
public void createGenericContainerExposesOpenHandles() {
Block<LLMBlockType> internalBlock = blocksController.create(LLMBlockConfiguration.builder()
.name("Analyze")
.llmDescriptor(LLMDescriptor.builder()
@ -92,42 +96,60 @@ public class ContainersControllerTest {
Container<GenericContainerType> container = containersController.create(GenericContainerConfiguration.builder()
.name("Container")
.subFlow(FlowData.builder().block(internalBlock).build())
.publicInputs(List.of(GenericContainerConfiguration.InputPort.builder()
.name("candidate")
.targetBlockId(internalBlock.getId())
.targetInputName("candidate")
.build()))
.publicOutputs(List.of(GenericContainerConfiguration.OutputPort.builder()
.name("analysis")
.sourceBlockId(internalBlock.getId())
.sourceOutputName("response")
.build()))
.build());
assertNotNull(container);
assertNotNull(container.getPosition());
assertEquals(0, container.getPosition().x());
assertEquals(0, container.getPosition().y());
assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate")));
assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("analysis")));
assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response")));
}
@Test
public void createGenericContainerAutoExposesOpenHandlesWhenPublicPortsAreMissing() {
Block<LLMBlockType> internalBlock = blocksController.create(LLMBlockConfiguration.builder()
.name("Analyze")
public void createGenericContainerPrefixesCollidingPortNamesWithNodeName() {
Block<LLMBlockType> firstInternalBlock = blocksController.create(LLMBlockConfiguration.builder()
.name("AnalyzeA")
.llmDescriptor(LLMDescriptor.builder()
.provider("testProvider")
.model("testModel")
.build())
.prompt("Analyze ${{candidate}}")
.build());
Block<LLMBlockType> secondInternalBlock = blocksController.create(LLMBlockConfiguration.builder()
.name("AnalyzeB")
.llmDescriptor(LLMDescriptor.builder()
.provider("testProvider")
.model("testModel")
.build())
.prompt("Review ${{candidate}}")
.build());
Container<GenericContainerType> container = containersController.create(GenericContainerConfiguration.builder()
.name("Container")
.subFlow(FlowData.builder().block(internalBlock).build())
.subFlow(FlowData.builder()
.block(Block.<LLMBlockType>builder()
.inputs(firstInternalBlock.getInputs())
.outputs(firstInternalBlock.getOutputs())
.specificConfiguration(firstInternalBlock.getSpecificConfiguration())
.type(firstInternalBlock.getType())
.position(new Position(0, 0))
.build())
.block(Block.<LLMBlockType>builder()
.inputs(secondInternalBlock.getInputs())
.outputs(secondInternalBlock.getOutputs())
.specificConfiguration(secondInternalBlock.getSpecificConfiguration())
.type(secondInternalBlock.getType())
.position(new Position(0, 0))
.build())
.build())
.build());
assertNotNull(container);
assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate")));
assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response")));
assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("AnalyzeA.candidate")));
assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("AnalyzeB.candidate")));
assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("AnalyzeA.response")));
assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("AnalyzeB.response")));
}
@Test

View File

@ -311,6 +311,62 @@ public class FlowControllerTest {
assertTrue(flows.stream().anyMatch(f -> publicOtherCreated.id().equals(f.id())));
}
@Test
public void getFlowByIdRespectsVisibilityRules() {
LoginEntity testUser = new LoginEntity("testuser", "testpassword");
LoginEntity otherUser = new LoginEntity("otheruser", "testpassword");
LLMDescriptor llmDescriptor = LLMDescriptor.builder()
.provider("testProvider")
.model("testModel")
.build();
Flow ownFlow = flowTestCreator.createFlowWithConnection(llmDescriptor);
FlowView ownCreated = flowController.createFlow(
new FlowCreateRequest(
ownFlow.getName(),
ownFlow.getDescription(),
FlowData.builder().blocks(ownFlow.getBlocks()).connections(ownFlow.getConnections()).build()),
testUser).getBody();
assertNotNull(ownCreated);
Flow publicOtherFlow = flowTestCreator.createFlowWithConnection(llmDescriptor);
FlowView publicOtherCreated = flowController.createFlow(
new FlowCreateRequest(
publicOtherFlow.getName(),
publicOtherFlow.getDescription(),
FlowData.builder().blocks(publicOtherFlow.getBlocks()).connections(publicOtherFlow.getConnections()).build()),
otherUser).getBody();
assertNotNull(publicOtherCreated);
flowRepository.findById(publicOtherCreated.id()).ifPresent(entity -> {
entity.setPublished(true);
flowRepository.save(entity);
});
Flow privateOtherFlow = flowTestCreator.createFlowWithConnection(llmDescriptor);
FlowView privateOtherCreated = flowController.createFlow(
new FlowCreateRequest(
privateOtherFlow.getName(),
privateOtherFlow.getDescription(),
FlowData.builder().blocks(privateOtherFlow.getBlocks()).connections(privateOtherFlow.getConnections()).build()),
otherUser).getBody();
assertNotNull(privateOtherCreated);
ResponseEntity<FlowView> ownResponse = flowController.getFlow(ownCreated.id(), testUser);
assertEquals(200, ownResponse.getStatusCode().value());
assertNotNull(ownResponse.getBody());
assertEquals(ownCreated.id(), ownResponse.getBody().id());
ResponseEntity<FlowView> publicResponse = flowController.getFlow(publicOtherCreated.id(), testUser);
assertEquals(200, publicResponse.getStatusCode().value());
assertNotNull(publicResponse.getBody());
assertEquals(publicOtherCreated.id(), publicResponse.getBody().id());
ResponseEntity<FlowView> forbiddenResponse = flowController.getFlow(privateOtherCreated.id(), testUser);
assertEquals(403, forbiddenResponse.getStatusCode().value());
}
@Test
public void createFlowRejectsBlockTamperedOutsideFactory() throws JsonProcessingException {
LLMDescriptor llmDescriptor = LLMDescriptor.builder()

View File

@ -341,22 +341,12 @@ public class ExecutionTest {
Container<GenericContainerType> container = genericContainerFactory.create(GenericContainerConfiguration.builder()
.name("Container")
.subFlow(FlowData.builder().block(internalBlock).build())
.publicInputs(List.of(GenericContainerConfiguration.InputPort.builder()
.name("candidateName")
.targetBlockId(internalBlock.getId())
.targetInputName("name")
.build()))
.publicOutputs(List.of(GenericContainerConfiguration.OutputPort.builder()
.name("analysis")
.sourceBlockId(internalBlock.getId())
.sourceOutputName("response")
.build()))
.build());
FlowData flow = FlowData.builder().container(container).build();
ExecutionObject execObject = executionsService.createExecution("Container flow", flow);
String exposedInputName = "candidateName";
String exposedOutputName = "analysis";
String exposedInputName = "name";
String exposedOutputName = "response";
executionsService.prepareInput(execObject.getId(), container.getId(), exposedInputName, "Frank");
execObject = executionsService.startExecution(execObject.getId());