From ad42c38259d6065b5ed70c7225fc99da4a33bbc7 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Thu, 19 Mar 2026 15:14:54 +0100 Subject: [PATCH] Add chat interaction runtime contract and schema metadata --- .../assistant/BlockCatalogService.java | 33 ++ .../assistant/FlowAssistantService.java | 1 + .../ChatInteractionBlockConfiguration.java | 69 +++ .../configurations/ChatInteractionInput.java | 36 ++ ...DynamicBlockConfigurationTypeResolver.java | 6 + .../configurations/JsonSchemaProducer.java | 467 ++++++++++++++++++ .../ChatInteractionBlockFactory.java | 116 +++++ .../manager/blocks/types/BlockTypes.java | 10 +- .../types/ChatInteractionBlockType.java | 37 ++ .../annotations/SchemaAllowedValues.java | 12 + .../annotations/UiDescription.java | 12 + .../annotations/UiEnabledWhen.java | 18 + .../configurations/annotations/UiLabel.java | 12 + .../annotations/UiOptionsFromNode.java | 16 + .../annotations/UiUniqueItemsBy.java | 12 + .../IteratorContainerConfiguration.java | 27 + .../factories/IteratorContainerFactory.java | 42 +- .../ContainerFlowInterfaceResolver.java | 23 +- .../IteratorContainerInterfaceResolver.java | 51 +- .../manager/controllers/BlocksController.java | 50 +- .../controllers/ContainersController.java | 43 +- .../manager/executions/ExecutionContext.java | 22 +- .../manager/executions/ExecutionListener.java | 2 + .../manager/executions/ExecutionsService.java | 4 + .../manager/executions/InteractionResult.java | 19 + .../executions/executors/NodeExecutors.java | 9 + .../executors/blocks/BlockExecutor.java | 6 + .../blocks/ChatInteractionExecutor.java | 168 +++++++ .../blocks/HumanInteractionExecutor.java | 11 + .../manager/executions/steps/Step.java | 23 +- .../flows/validation/FlowDataValidator.java | 35 +- .../validation/FlowExecutionValidator.java | 27 + .../workflow/manager/llms/ChatMessage.java | 17 + .../manager/llms/providers/LLMProvider.java | 19 + .../providers/google/GeminiLLMProvider.java | 31 +- .../ollama/InternalOllamaLLMProvider.java | 58 +++ .../ollama/response/ChatResponse.java | 13 + .../ollama/response/ChatResponseMessage.java | 14 + .../resources/workflow-editor-init/flows.json | 315 ++++++++++++ .../controllers/BlocksControllerTest.java | 214 ++++++++ .../controllers/ContainersControllerTest.java | 98 +++- .../controllers/ExecutionControllerTest.java | 6 +- .../controllers/FlowControllerTest.java | 75 ++- .../manager/executions/ExecutionTest.java | 98 +++- .../executions/ExecutionWithContainer.java | 24 +- 45 files changed, 2231 insertions(+), 170 deletions(-) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionBlockConfiguration.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionInput.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/factories/ChatInteractionBlockFactory.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/types/ChatInteractionBlockType.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/SchemaAllowedValues.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiDescription.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiEnabledWhen.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiLabel.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiOptionsFromNode.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiUniqueItemsBy.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/InteractionResult.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/llms/ChatMessage.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponse.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponseMessage.java diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockCatalogService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockCatalogService.java index cccbc91..075015f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockCatalogService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockCatalogService.java @@ -25,11 +25,21 @@ public class BlockCatalogService { String type, String description, boolean userInteractive, + AssistantInteractionContractDescriptor interactionContract, String configurationClass, JsonNode schema, Block exampleBlock) { } + public record AssistantInteractionContractDescriptor( + String kind, + String messageField, + String completionField, + String historyField, + String responseField, + boolean supportsPartialResult) { + } + public record AssistantPromptFieldDescriptor( String name, String type, @@ -87,6 +97,7 @@ public class BlockCatalogService { blockType.getName(), blockType.getDescription(), blockType.isUserInteractive(), + resolveInteractionContract(blockType), configurationClass == null ? null : configurationClass.getName(), schema, exampleBlock); @@ -176,4 +187,26 @@ public class BlockCatalogService { fields.forEachRemaining(entries::add); return entries; } + + private AssistantInteractionContractDescriptor resolveInteractionContract(BlockType blockType) { + if ("ChatInteraction".equals(blockType.getName())) { + return new AssistantInteractionContractDescriptor( + "chat-session", + "message", + "response", + "history", + "response", + true); + } + if ("HumanInteractionBlock".equals(blockType.getName())) { + return new AssistantInteractionContractDescriptor( + "single-response", + null, + "output", + null, + "output", + false); + } + return null; + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java index 955e111..bfda445 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java @@ -287,6 +287,7 @@ public class FlowAssistantService { switch (descriptor.configurationType()) { case "LLMBlockConfiguration" -> config.set("llmDescriptor", llmDescriptorNode(model)); case "HumanInteractiveBlockConfiguration" -> config.set("simulateWith", llmDescriptorNode(model)); + case "ChatInteractionBlockConfiguration" -> config.set("llmDescriptor", llmDescriptorNode(model)); case "ConditionalBlockConfiguration" -> { boolean useLlm = inferConditionalUseLlm(config); config.put("useLlm", useLlm); diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionBlockConfiguration.java new file mode 100644 index 0000000..f3170fd --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionBlockConfiguration.java @@ -0,0 +1,69 @@ +package it.cnr.isti.workflow.manager.blocks.configurations; + +import java.util.List; + +import com.fasterxml.jackson.annotation.JsonProperty; + +import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; +import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiUniqueItemsBy; +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import jakarta.validation.Valid; +import jakarta.validation.constraints.AssertTrue; +import jakarta.validation.constraints.NotNull; +import lombok.Builder; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import lombok.NonNull; + +@Data +@EqualsAndHashCode(callSuper = true) +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +public class ChatInteractionBlockConfiguration extends BlockConfiguration { + + @NotNull + @Valid + @JsonProperty(required = true) + private LLMDescriptor llmDescriptor; + + @Structural + @UiUniqueItemsBy("name") + @Valid + private List inputs = List.of(); + + @Override + public Class getBlockType() { + return ChatInteractionBlockType.class; + } + + @Builder + public ChatInteractionBlockConfiguration(@NonNull String name, LLMDescriptor llmDescriptor, + List inputs) { + super(name); + this.llmDescriptor = llmDescriptor; + this.inputs = inputs == null ? List.of() : List.copyOf(inputs); + } + + public static ChatInteractionBlockConfiguration empty() { + ChatInteractionBlockConfiguration configuration = new ChatInteractionBlockConfiguration(); + configuration.name = ChatInteractionBlockType.TYPE; + configuration.inputs = List.of(); + return configuration; + } + + @AssertTrue(message = "inputs must have unique names") + boolean areInputNamesUnique() { + if (inputs == null || inputs.isEmpty()) { + return true; + } + return inputs.stream() + .map(ChatInteractionInput::name) + .filter(name -> name != null && !name.isBlank()) + .distinct() + .count() == inputs.stream() + .map(ChatInteractionInput::name) + .filter(name -> name != null && !name.isBlank()) + .count(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionInput.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionInput.java new file mode 100644 index 0000000..5efb6a0 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/ChatInteractionInput.java @@ -0,0 +1,36 @@ +package it.cnr.isti.workflow.manager.blocks.configurations; + +import com.fasterxml.jackson.annotation.JsonAlias; +import com.fasterxml.jackson.annotation.JsonProperty; + +import it.cnr.isti.workflow.manager.configurations.annotations.SchemaAllowedValues; +import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; +import it.cnr.isti.workflow.manager.ios.IOType; +import jakarta.validation.constraints.AssertTrue; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; +import jakarta.validation.constraints.Size; + +public record ChatInteractionInput( + @Size(max = 8) + @NotBlank String name, + @UiLabel("type") + @SchemaAllowedValues({ "TEXT" }) + @JsonProperty("ioType") + @JsonAlias("type") + @NotNull IOType ioType, + boolean multiple) { + + public ChatInteractionInput { + ioType = ioType == null ? IOType.TEXT : ioType; + } + + public ChatInteractionInput(String name, IOType ioType) { + this(name, ioType, false); + } + + @AssertTrue(message = "ChatInteraction inputs support only TEXT or TEXT[]") + boolean hasSupportedType() { + return ioType == IOType.TEXT; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java index f6d9b6b..b5c4b00 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/DynamicBlockConfigurationTypeResolver.java @@ -14,6 +14,9 @@ import java.util.Map; public class DynamicBlockConfigurationTypeResolver extends TypeIdResolverBase { + private static final String CHAT_INTERACTION_LEGACY_CONFIGURATION_ID = "ChatHumanInteractionBlockConfiguration"; + private static final String CHAT_INTERACTION_CONFIGURATION_ID = "ChatInteractionBlockConfiguration"; + private Map> idToClass = new HashMap<>(); private Map, String> classToId = new HashMap<>(); @@ -44,6 +47,9 @@ public class DynamicBlockConfigurationTypeResolver extends TypeIdResolverBase { @Override public JavaType typeFromId(DatabindContext context, String id) { + if (CHAT_INTERACTION_LEGACY_CONFIGURATION_ID.equals(id)) { + id = CHAT_INTERACTION_CONFIGURATION_ID; + } Class clazz = idToClass.get(id); if (clazz == null) { throw new IllegalArgumentException("Unknown BlockConfiguration type id: " + id); diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java index 89e92bf..2b81e69 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java @@ -28,7 +28,14 @@ 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.SchemaAllowedValues; +import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription; import it.cnr.isti.workflow.manager.configurations.annotations.UiRequiredWhen; +import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen; +import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel; +import it.cnr.isti.workflow.manager.configurations.annotations.UiOptionsFromNode; +import it.cnr.isti.workflow.manager.configurations.annotations.UiUniqueItemsBy; +import jakarta.validation.constraints.Size; @Component public class JsonSchemaProducer { @@ -50,14 +57,28 @@ public class JsonSchemaProducer { Map, Map> longTextMap = collectLongTextMetadata(type); Map, Map> structuralMap = collectStructuralMetadata(type); Map, Map> uiDependencyMap = collectUiDependencyMetadata(type); + Map, Map> uiEnabledWhenMap = collectUiEnabledWhenMetadata(type); + Map, Map> uiOptionsFromNodeMap = collectUiOptionsFromNodeMetadata(type); Map, Map> uiRequiredWhenMap = collectUiRequiredWhenMetadata(type); + Map, Map> uiUniqueItemsByMap = collectUiUniqueItemsByMetadata(type); + Map, Map> uiLabelMap = collectUiLabelMetadata(type); + Map, Map> uiDescriptionMap = collectUiDescriptionMetadata(type); + Map, Map> sizeMap = collectSizeMetadata(type); + Map, Map> schemaAllowedValuesMap = collectSchemaAllowedValuesMetadata(type); Map, Map> 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())); + applyUiEnabledWhenMetadata(root, uiEnabledWhenMap.getOrDefault(type, Map.of())); + applyUiOptionsFromNodeMetadata(root, uiOptionsFromNodeMap.getOrDefault(type, Map.of())); applyUiRequiredWhenMetadata(root, uiRequiredWhenMap.getOrDefault(type, Map.of())); + applyUiUniqueItemsByMetadata(root, uiUniqueItemsByMap.getOrDefault(type, Map.of())); + applyUiLabelMetadata(root, uiLabelMap.getOrDefault(type, Map.of())); + applyUiDescriptionMetadata(root, uiDescriptionMap.getOrDefault(type, Map.of())); + applySizeMetadata(root, sizeMap.getOrDefault(type, Map.of())); + applySchemaAllowedValuesMetadata(root, schemaAllowedValuesMap.getOrDefault(type, Map.of())); applyConfigurableAsInputMetadata(root, configurableAsInputMap.getOrDefault(type, Map.of())); JsonNode definitionsNode = root.has("definitions") ? root.get("definitions") : root.get("$defs"); @@ -67,7 +88,14 @@ public class JsonSchemaProducer { metadataClasses.addAll(longTextMap.keySet()); metadataClasses.addAll(structuralMap.keySet()); metadataClasses.addAll(uiDependencyMap.keySet()); + metadataClasses.addAll(uiEnabledWhenMap.keySet()); + metadataClasses.addAll(uiOptionsFromNodeMap.keySet()); metadataClasses.addAll(uiRequiredWhenMap.keySet()); + metadataClasses.addAll(uiUniqueItemsByMap.keySet()); + metadataClasses.addAll(uiLabelMap.keySet()); + metadataClasses.addAll(uiDescriptionMap.keySet()); + metadataClasses.addAll(sizeMap.keySet()); + metadataClasses.addAll(schemaAllowedValuesMap.keySet()); metadataClasses.addAll(configurableAsInputMap.keySet()); for (Entry entry : iterable(definitions.fields())) { if (!(entry.getValue() instanceof ObjectNode classSchema)) { @@ -80,7 +108,14 @@ public class JsonSchemaProducer { applyLongTextMetadata(classSchema, longTextMap.get(matchedClass)); applyStructuralMetadata(classSchema, structuralMap.get(matchedClass)); applyUiDependencyMetadata(classSchema, uiDependencyMap.get(matchedClass)); + applyUiEnabledWhenMetadata(classSchema, uiEnabledWhenMap.get(matchedClass)); + applyUiOptionsFromNodeMetadata(classSchema, uiOptionsFromNodeMap.get(matchedClass)); applyUiRequiredWhenMetadata(classSchema, uiRequiredWhenMap.get(matchedClass)); + applyUiUniqueItemsByMetadata(classSchema, uiUniqueItemsByMap.get(matchedClass)); + applyUiLabelMetadata(classSchema, uiLabelMap.get(matchedClass)); + applyUiDescriptionMetadata(classSchema, uiDescriptionMap.get(matchedClass)); + applySizeMetadata(classSchema, sizeMap.get(matchedClass)); + applySchemaAllowedValuesMetadata(classSchema, schemaAllowedValuesMap.get(matchedClass)); applyConfigurableAsInputMetadata(classSchema, configurableAsInputMap.get(matchedClass)); } } @@ -253,6 +288,162 @@ public class JsonSchemaProducer { return result; } + private Map, Map> collectUiLabelMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + UiLabel uiLabel = field.getAnnotation(UiLabel.class); + if (uiLabel != null) { + metadata.put(field.getName(), uiLabel); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + UiLabel uiLabel = component.getAnnotation(UiLabel.class); + if (uiLabel != null) { + metadata.put(component.getName(), uiLabel); + } + enqueueRelatedTypes(queue, component.getGenericType(), component.getType()); + } + } + + if (!metadata.isEmpty()) { + result.put(current, metadata); + } + } + + return result; + } + + private Map, Map> collectUiDescriptionMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + UiDescription uiDescription = field.getAnnotation(UiDescription.class); + if (uiDescription != null) { + metadata.put(field.getName(), uiDescription); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + UiDescription uiDescription = component.getAnnotation(UiDescription.class); + if (uiDescription != null) { + metadata.put(component.getName(), uiDescription); + } + enqueueRelatedTypes(queue, component.getGenericType(), component.getType()); + } + } + + if (!metadata.isEmpty()) { + result.put(current, metadata); + } + } + + return result; + } + + private Map, Map> collectSizeMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + Size size = field.getAnnotation(Size.class); + if (size != null) { + metadata.put(field.getName(), size); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + Size size = component.getAnnotation(Size.class); + if (size != null) { + metadata.put(component.getName(), size); + } + enqueueRelatedTypes(queue, component.getGenericType(), component.getType()); + } + } + + if (!metadata.isEmpty()) { + result.put(current, metadata); + } + } + + return result; + } + + private Map, Map> collectSchemaAllowedValuesMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + SchemaAllowedValues allowedValues = field.getAnnotation(SchemaAllowedValues.class); + if (allowedValues != null) { + metadata.put(field.getName(), allowedValues); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + SchemaAllowedValues allowedValues = component.getAnnotation(SchemaAllowedValues.class); + if (allowedValues != null) { + metadata.put(component.getName(), allowedValues); + } + enqueueRelatedTypes(queue, component.getGenericType(), component.getType()); + } + } + + if (!metadata.isEmpty()) { + result.put(current, metadata); + } + } + + return result; + } + private void applyDynamicSchemaMetadata(ObjectNode classSchema, Map metadata) { if (metadata == null || metadata.isEmpty()) { return; @@ -404,6 +595,45 @@ public class JsonSchemaProducer { return result; } + private Map, Map> collectUiUniqueItemsByMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + UiUniqueItemsBy annotation = field.getAnnotation(UiUniqueItemsBy.class); + if (annotation != null) { + metadata.put(field.getName(), annotation); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + UiUniqueItemsBy annotation = component.getAnnotation(UiUniqueItemsBy.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 applyUiDependencyMetadata(ObjectNode classSchema, Map metadata) { if (metadata == null || metadata.isEmpty()) { return; @@ -437,6 +667,109 @@ public class JsonSchemaProducer { } } + private void applyUiUniqueItemsByMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + propertySchema.put("x-ui-unique-by", entry.getValue().value()); + } + } + + private void applyUiLabelMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + if (!entry.getValue().value().isBlank()) { + propertySchema.put("x-ui-label", entry.getValue().value()); + } + } + } + + private void applyUiDescriptionMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + if (!entry.getValue().value().isBlank()) { + propertySchema.put("x-ui-description", entry.getValue().value()); + } + } + } + + private void applySizeMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + Size size = entry.getValue(); + if (size.min() > 0) { + propertySchema.put("minLength", size.min()); + } + if (size.max() < Integer.MAX_VALUE) { + propertySchema.put("maxLength", size.max()); + } + } + } + + private void applySchemaAllowedValuesMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + ArrayNode enumValues = propertySchema.putArray("enum"); + for (String value : entry.getValue().value()) { + enumValues.add(value); + } + } + } + private Map, Map> collectUiRequiredWhenMetadata(Class rootClass) { Map, Map> result = new HashMap<>(); Set> visited = new HashSet<>(); @@ -476,6 +809,140 @@ public class JsonSchemaProducer { return result; } + private Map, Map> collectUiOptionsFromNodeMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + UiOptionsFromNode annotation = field.getAnnotation(UiOptionsFromNode.class); + if (annotation != null) { + metadata.put(field.getName(), annotation); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + UiOptionsFromNode annotation = component.getAnnotation(UiOptionsFromNode.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 applyUiOptionsFromNodeMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + + UiOptionsFromNode annotation = entry.getValue(); + ObjectNode optionsFromNode = propertySchema.putObject("x-ui-options-from-node"); + optionsFromNode.put("collection", annotation.collection()); + optionsFromNode.put("valueField", annotation.valueField()); + optionsFromNode.put("labelField", annotation.labelField()); + } + } + + private Map, Map> collectUiEnabledWhenMetadata(Class rootClass) { + Map, Map> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + + Map metadata = new LinkedHashMap<>(); + for (Field field : current.getDeclaredFields()) { + UiEnabledWhen annotation = field.getAnnotation(UiEnabledWhen.class); + if (annotation != null) { + metadata.put(field.getName(), annotation); + } + enqueueRelatedTypes(queue, field.getGenericType(), field.getType()); + } + + if (current.isRecord()) { + for (RecordComponent component : current.getRecordComponents()) { + UiEnabledWhen annotation = component.getAnnotation(UiEnabledWhen.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 applyUiEnabledWhenMetadata(ObjectNode classSchema, Map metadata) { + if (metadata == null || metadata.isEmpty()) { + return; + } + JsonNode propsNode = classSchema.get("properties"); + if (!(propsNode instanceof ObjectNode properties)) { + return; + } + + for (Entry entry : metadata.entrySet()) { + JsonNode propNode = properties.get(entry.getKey()); + if (!(propNode instanceof ObjectNode propertySchema)) { + continue; + } + + UiEnabledWhen dependency = entry.getValue(); + ObjectNode enabledWhen = propertySchema.putObject("x-ui-enabled-when"); + enabledWhen.put("field", dependency.field()); + if (!dependency.equals().isBlank()) { + enabledWhen.put("equals", dependency.equals()); + } + if (dependency.equalsAny().length > 0) { + ArrayNode equalsAny = enabledWhen.putArray("in"); + for (String value : dependency.equalsAny()) { + equalsAny.add(value); + } + } + if (dependency.present()) { + enabledWhen.put("present", true); + } + } + } + private void applyUiRequiredWhenMetadata(ObjectNode classSchema, Map metadata) { if (metadata == null || metadata.isEmpty()) { return; diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/ChatInteractionBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/ChatInteractionBlockFactory.java new file mode 100644 index 0000000..84f80e1 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/ChatInteractionBlockFactory.java @@ -0,0 +1,116 @@ +package it.cnr.isti.workflow.manager.blocks.factories; + +import java.util.List; +import java.util.Objects; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.IOCapability; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionInput; +import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; +import it.cnr.isti.workflow.manager.ios.IODescriptor; +import it.cnr.isti.workflow.manager.ios.IOType; + +@Component +public class ChatInteractionBlockFactory + implements BlockFactory { + + public static final String INTERACTION_FIELD = "message"; + public static final String FINAL_RESPONSE_FIELD = "response"; + public static final String RESPONSE_OUTPUT = "response"; + public static final String HISTORY_OUTPUT = "history"; + + private static final List RESPONSE_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.TEXT, false)); + private static final List HISTORY_CAPABILITIES = List.of( + new IOCapability(IOCapabilityType.TEXT, true)); + + @Autowired + private ChatInteractionBlockType blockType; + + @Override + public Block create(ChatInteractionBlockConfiguration configuration) { + validateConfiguration(configuration); + return Block.builder() + .inputs(resolveInputs(configuration)) + .output(IODescriptor.output(RESPONSE_OUTPUT, IOType.TEXT, false, RESPONSE_CAPABILITIES)) + .output(IODescriptor.output(HISTORY_OUTPUT, IOType.TEXT, true, HISTORY_CAPABILITIES)) + .specificConfiguration(configuration) + .type(blockType) + .build(); + } + + @Override + public Block createEmpty() { + return create(ChatInteractionBlockConfiguration.empty()); + } + + @Override + public Class getBlockType() { + return ChatInteractionBlockType.class; + } + + @Override + public List supportedInputCapabilities() { + return List.of( + new IOCapability(IOCapabilityType.TEXT, false), + new IOCapability(IOCapabilityType.TEXT, true)); + } + + @Override + public List supportedOutputCapabilities() { + return List.of( + new IOCapability(IOCapabilityType.TEXT, false), + new IOCapability(IOCapabilityType.TEXT, true)); + } + + private List resolveInputs(ChatInteractionBlockConfiguration configuration) { + if (configuration.getInputs() == null || configuration.getInputs().isEmpty()) { + return List.of(); + } + return configuration.getInputs().stream() + .map(this::toDescriptor) + .toList(); + } + + private IODescriptor toDescriptor(ChatInteractionInput input) { + IOType type = input.ioType() == null ? IOType.TEXT : input.ioType(); + return IODescriptor.input(input.name(), type, input.multiple(), + List.of(new IOCapability(toCapabilityType(type), input.multiple()))); + } + + private void validateConfiguration(ChatInteractionBlockConfiguration configuration) { + if (configuration == null || configuration.getInputs() == null || configuration.getInputs().isEmpty()) { + return; + } + long distinctNames = configuration.getInputs().stream() + .map(ChatInteractionInput::name) + .filter(Objects::nonNull) + .distinct() + .count(); + long names = configuration.getInputs().stream() + .map(ChatInteractionInput::name) + .filter(Objects::nonNull) + .count(); + if (distinctNames != names) { + throw new IllegalArgumentException("inputs must have unique names"); + } + boolean unsupportedType = configuration.getInputs().stream() + .anyMatch(input -> input.ioType() != null && input.ioType() != IOType.TEXT); + if (unsupportedType) { + throw new IllegalArgumentException("ChatInteraction inputs support only TEXT or TEXT[]"); + } + } + + private IOCapabilityType toCapabilityType(IOType type) { + return switch (type) { + case FILE, CSV -> IOCapabilityType.FILE; + case TEXT -> IOCapabilityType.TEXT; + case ANY -> IOCapabilityType.ANY; + }; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockTypes.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockTypes.java index ff34ab8..0e9f0c7 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockTypes.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/BlockTypes.java @@ -10,10 +10,18 @@ import org.springframework.stereotype.Component; public class BlockTypes { private static Map blockTypes = new HashMap<>(); + private static final String CHAT_INTERACTION_LEGACY_TYPE = "ChatHumanInteraction"; public static BlockType get(String name) { - return blockTypes.get(name); + BlockType blockType = blockTypes.get(name); + if (blockType != null) { + return blockType; + } + if (CHAT_INTERACTION_LEGACY_TYPE.equals(name)) { + return blockTypes.get("ChatInteraction"); + } + return null; } BlockTypes(List types) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/ChatInteractionBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/ChatInteractionBlockType.java new file mode 100644 index 0000000..7ceb102 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/ChatInteractionBlockType.java @@ -0,0 +1,37 @@ +package it.cnr.isti.workflow.manager.blocks.types; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration; + +@Component(ChatInteractionBlockType.TYPE) +public class ChatInteractionBlockType implements BlockType { + + public static final String TYPE = "ChatInteraction"; + + @Override + public String getName() { + return TYPE; + } + + @Override + public String getDescription() { + return "A human-interactive chat block backed by an LLM with persistent conversation history"; + } + + @Override + public boolean validate() { + return true; + } + + @Override + public boolean isUserInteractive() { + return true; + } + + @Override + public Class> getBlockConfigurationClass() { + return ChatInteractionBlockConfiguration.class; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/SchemaAllowedValues.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/SchemaAllowedValues.java new file mode 100644 index 0000000..f21f5d7 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/SchemaAllowedValues.java @@ -0,0 +1,12 @@ +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 SchemaAllowedValues { + String[] value(); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiDescription.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiDescription.java new file mode 100644 index 0000000..8fb0e5e --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiDescription.java @@ -0,0 +1,12 @@ +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 UiDescription { + String value(); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiEnabledWhen.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiEnabledWhen.java new file mode 100644 index 0000000..3732ce8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiEnabledWhen.java @@ -0,0 +1,18 @@ +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 UiEnabledWhen { + String field(); + + String equals() default ""; + + String[] equalsAny() default {}; + + boolean present() default false; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiLabel.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiLabel.java new file mode 100644 index 0000000..d2ec2cb --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiLabel.java @@ -0,0 +1,12 @@ +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 UiLabel { + String value(); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiOptionsFromNode.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiOptionsFromNode.java new file mode 100644 index 0000000..85ec091 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiOptionsFromNode.java @@ -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 UiOptionsFromNode { + String collection(); + + String valueField() default "name"; + + String labelField() default "name"; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiUniqueItemsBy.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiUniqueItemsBy.java new file mode 100644 index 0000000..d410c1a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiUniqueItemsBy.java @@ -0,0 +1,12 @@ +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 UiUniqueItemsBy { + String value(); +} 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 index 3653fad..5950b58 100644 --- 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 @@ -1,12 +1,17 @@ package it.cnr.isti.workflow.manager.containers.configurations; import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonIgnore; import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever; +import it.cnr.isti.workflow.manager.configurations.annotations.UiOptionsFromNode; import it.cnr.isti.workflow.manager.configurations.annotations.Structural; +import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; import it.cnr.isti.workflow.manager.flows.model.FlowData; import jakarta.validation.Valid; +import jakarta.validation.constraints.AssertTrue; import lombok.Builder; import lombok.EqualsAndHashCode; import lombok.Getter; @@ -30,6 +35,8 @@ public class IteratorContainerConfiguration extends ContainerConfiguration handle.publicName().equals(iterationInput) && !handle.handle().io().isMultiple()); + } + + private boolean hasSubFlowNodes() { + return subFlow != null && !subFlow.getNodes().isEmpty(); + } } 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 index d3f5267..70a23f4 100644 --- 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 @@ -6,6 +6,7 @@ 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.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; @@ -22,13 +23,12 @@ public class IteratorContainerFactory implements ContainerFactory create(IteratorContainerConfiguration configuration) { - ContainerFlowInterface exposedInterface = isEmpty(configuration) - ? new ContainerFlowInterface(java.util.List.of(), java.util.List.of()) - : IteratorContainerInterfaceResolver.resolve(configuration); + IteratorContainerConfiguration effectiveConfiguration = applyDefaultIterationInput(configuration); + ContainerFlowInterface exposedInterface = resolveInterface(effectiveConfiguration); return Container.builder() .inputs(exposedInterface.inputs()) .outputs(exposedInterface.outputs()) - .specificConfiguration(configuration) + .specificConfiguration(effectiveConfiguration) .type(containerType) .position(DEFAULT_POSITION) .build(); @@ -49,4 +49,38 @@ public class IteratorContainerFactory implements ContainerFactory openInputs = + ContainerFlowInterfaceResolver.getExposedInputs(configuration.getSubFlow()); + if (openInputs.size() != 1) { + return configuration; + } + return IteratorContainerConfiguration.builder() + .name(configuration.getName()) + .subFlow(configuration.getSubFlow()) + .iterationInput(openInputs.getFirst().publicName()) + .build(); + } + + private ContainerFlowInterface resolveInterface(IteratorContainerConfiguration configuration) { + if (isEmpty(configuration)) { + return new ContainerFlowInterface(java.util.List.of(), java.util.List.of()); + } + if (configuration.getIterationInput() == null || configuration.getIterationInput().isBlank()) { + return new ContainerFlowInterface( + ContainerFlowInterfaceResolver.resolveGenericInputs(configuration.getSubFlow()), + ContainerFlowInterfaceResolver.resolveGenericOutputs(configuration.getSubFlow())); + } + try { + return IteratorContainerInterfaceResolver.resolve(configuration); + } catch (IllegalArgumentException exception) { + return new ContainerFlowInterface( + ContainerFlowInterfaceResolver.resolveGenericInputs(configuration.getSubFlow()), + ContainerFlowInterfaceResolver.resolveGenericOutputs(configuration.getSubFlow())); + } + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java index fb74e0e..4d285ed 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/containers/iresolvers/ContainerFlowInterfaceResolver.java @@ -33,16 +33,21 @@ public final class ContainerFlowInterfaceResolver { public static ContainerFlowInterface resolve(GenericContainerConfiguration configuration) { FlowData subFlow = configuration == null ? null : configuration.getSubFlow(); - List openInputs = getExposedInputs(subFlow); - List openOutputs = getExposedOutputs(subFlow); - List inputs = openInputs.stream() - .map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io())) - .toList(); - List outputs = openOutputs.stream() - .map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io())) - .toList(); + return new ContainerFlowInterface(resolveGenericInputs(subFlow), resolveGenericOutputs(subFlow)); + } - return new ContainerFlowInterface(inputs, outputs); + public static List resolveGenericInputs(FlowData subFlow) { + List openInputs = getExposedInputs(subFlow); + return openInputs.stream() + .map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io())) + .toList(); + } + + public static List resolveGenericOutputs(FlowData subFlow) { + List openOutputs = getExposedOutputs(subFlow); + return openOutputs.stream() + .map(handle -> cloneDescriptor(handle.publicName(), handle.handle().io())) + .toList(); } public static List getExposedInputs(FlowData 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 index 470921c..6f44820 100644 --- 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 @@ -2,10 +2,8 @@ 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; @@ -49,11 +47,7 @@ public final class IteratorContainerInterfaceResolver { "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()); + String iteratedPublicName = iteratedHandle.publicName(); List resolvedInputs = new ArrayList<>(); for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedInputs) { @@ -91,49 +85,6 @@ public final class IteratorContainerInterfaceResolver { 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, diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java index 00fe90b..826ac68 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/BlocksController.java @@ -29,6 +29,7 @@ import com.fasterxml.jackson.core.JsonProcessingException; public class BlocksController { private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(BlocksController.class); + private static final String CHAT_INTERACTION_LEGACY_TYPE = "ChatHumanInteraction"; @Autowired Map blockTypes; @@ -43,6 +44,7 @@ public class BlocksController { String type, String description, boolean userInteractive, + InteractionContractDescriptor interactionContract, boolean hasExampleBlock, String exampleBlockEndpoint, String configurationType, @@ -50,6 +52,15 @@ public class BlocksController { Object schema) { } + public record InteractionContractDescriptor( + String kind, + String messageField, + String completionField, + String historyField, + String responseField, + boolean supportsPartialResult) { + } + @GetMapping("types") @Operation(summary = "Get block types", description = "Returns all block types with their configuration descriptor and JSON schema.") @@ -66,7 +77,7 @@ public class BlocksController { @GetMapping("/types/{type}/configuration/descriptor") @Operation(summary = "Get block configuration descriptor by type", description = "Returns the descriptor and JSON schema for the requested block type.") public BlockConfigurationDescriptor getConfigurationDescriptorForType(@PathVariable String type) { - BlockType blockType = blockTypes.get(type); + BlockType blockType = resolveBlockType(type); if (blockType == null) { throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Block type not found: " + type); } @@ -77,7 +88,7 @@ public class BlocksController { @GetMapping("/types/{type}/example") @Operation(summary = "Get block example by type", description = "Returns an empty example block for the requested block type, intended for UI scaffolding.") public Block getExampleForType(@PathVariable String type) { - BlockType blockType = blockTypes.get(type); + BlockType blockType = resolveBlockType(type); if (blockType == null) { throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Block type not found: " + type); } @@ -97,6 +108,7 @@ public class BlocksController { blockType.getName(), blockType.getDescription(), blockType.isUserInteractive(), + resolveInteractionContract(blockType), true, getExampleEndpoint(blockType), null, @@ -109,6 +121,7 @@ public class BlocksController { blockType.getName(), blockType.getDescription(), blockType.isUserInteractive(), + resolveInteractionContract(blockType), true, getExampleEndpoint(blockType), configurationClass.getSimpleName(), @@ -121,6 +134,39 @@ public class BlocksController { return "/blocks/types/" + blockType.getName() + "/example"; } + private InteractionContractDescriptor resolveInteractionContract(BlockType blockType) { + if ("ChatInteraction".equals(blockType.getName())) { + return new InteractionContractDescriptor( + "chat-session", + "message", + "response", + "history", + "response", + true); + } + if ("HumanInteractionBlock".equals(blockType.getName())) { + return new InteractionContractDescriptor( + "single-response", + null, + "output", + null, + "output", + false); + } + return null; + } + + private BlockType resolveBlockType(String type) { + BlockType blockType = blockTypes.get(type); + if (blockType != null) { + return blockType; + } + if (CHAT_INTERACTION_LEGACY_TYPE.equals(type)) { + return blockTypes.get("ChatInteraction"); + } + return null; + } + @SuppressWarnings("unchecked") @PostMapping @SecurityRequirement(name = "bearerAuth") 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 8e5a29c..653e2ee 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 @@ -2,11 +2,8 @@ package it.cnr.isti.workflow.manager.controllers; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; - import org.eclipse.microprofile.openapi.annotations.Operation; import org.eclipse.microprofile.openapi.annotations.security.SecurityRequirement; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; @@ -31,17 +28,24 @@ import it.cnr.isti.workflow.manager.ios.IODescriptor; @RequestMapping("/containers") public class ContainersController { - @Autowired Map containerTypes; - @Autowired List> containerFactories; - @Autowired JsonSchemaProducer schemaProducer; - @Autowired - ContainerSubFlowValidator containerSubFlowValidator; + private final ContainerSubFlowValidator containerSubFlowValidator; + + public ContainersController( + Map containerTypes, + List> containerFactories, + JsonSchemaProducer schemaProducer, + ContainerSubFlowValidator containerSubFlowValidator) { + this.containerTypes = containerTypes; + this.containerFactories = containerFactories; + this.schemaProducer = schemaProducer; + this.containerSubFlowValidator = containerSubFlowValidator; + } public record ContainerConfigurationDescriptor( String type, @@ -115,7 +119,6 @@ public class ContainersController { .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); } @@ -135,28 +138,6 @@ public class ContainersController { private ContainerHandleDescriptor toHandle(ContainerFlowInterfaceResolver.OpenHandle handle) { return new ContainerHandleDescriptor(handle.blockId(), handle.blockName(), handle.io()); } - - private > void validateConfiguration(C configuration, - ContainerFactory factory) { - if (configuration == null) { - return; - } - - ContainerSubFlowValidator.ValidationResult validation = containerSubFlowValidator - .validate(configuration.getSubFlow()); - 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); - } - } - private ContainerConfigurationDescriptor toDescriptor(ContainerType containerType) { Class> configurationClass = containerType.getContainerConfigurationClass(); return new ContainerConfigurationDescriptor( diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java index fee7577..4b7fd09 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java @@ -26,6 +26,9 @@ public class ExecutionContext implements ExecutionListener { @Getter() private Map result = new HashMap<>(); + @Getter() + private Map partialResult = new HashMap<>(); + Long startTime = null; Long endTime = null; @@ -59,6 +62,14 @@ public class ExecutionContext implements ExecutionListener { this.result.put(new FieldKey(nodeId, key), value); } + private void addPartialResult(String nodeId, String key, Object value) { + this.partialResult.put(new FieldKey(nodeId, key), value); + } + + private void clearPartialResults(String nodeId) { + this.partialResult.keySet().removeIf(key -> key.nodeId().equals(nodeId)); + } + private void addInput(String nodeId, String key, Object value) { this.inputs.put(new FieldKey(nodeId, key), value); } @@ -74,6 +85,7 @@ public class ExecutionContext implements ExecutionListener { @Override public void completed(String id, Map result) { logger.info("Step " + id + " completed with result: " + result); + clearPartialResults(id); Step completedStep = this.steps.get(id); completedStep.getOutputs().forEach(output -> { String outputName = output.getDescriptor().getName(); @@ -113,7 +125,9 @@ public class ExecutionContext implements ExecutionListener { public synchronized void paused(String id) { logger.info("Step " + id + " paused"); synchronized (this.waitingSteps) { - this.waitingSteps.add(id); + if (!this.waitingSteps.contains(id)) { + this.waitingSteps.add(id); + } if (this.status != ExecutionStatus.WAITING) { this.setStatus(ExecutionStatus.WAITING); } @@ -131,6 +145,12 @@ public class ExecutionContext implements ExecutionListener { } } + @Override + public synchronized void partialUpdated(String id, Map partialResult) { + clearPartialResults(id); + partialResult.forEach((key, value) -> addPartialResult(id, key, value)); + } + protected void setInput(String stepId, String inputName, Object value) { Step step = this.steps.get(stepId); if (step == null) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java index 53bc4fe..4fbf8a9 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionListener.java @@ -15,4 +15,6 @@ public interface ExecutionListener { void paused(String id); void resumed(String id); + + void partialUpdated(String id, Map partialResult); } 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 9bab894..45b05f5 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 @@ -11,6 +11,7 @@ import org.springframework.stereotype.Service; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; @@ -147,6 +148,9 @@ public class ExecutionsService { if (configuration instanceof HumanInteractiveBlockConfiguration humanConfiguration) { return humanConfiguration.getSimulateWith(); } + if (configuration instanceof ChatInteractionBlockConfiguration chatConfiguration) { + return chatConfiguration.getLlmDescriptor(); + } if (configuration instanceof ConditionalBlockConfiguration conditionalConfiguration && conditionalConfiguration.isUseLlm()) { return conditionalConfiguration.getLlmDescriptor(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/InteractionResult.java b/src/main/java/it/cnr/isti/workflow/manager/executions/InteractionResult.java new file mode 100644 index 0000000..261eb83 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/InteractionResult.java @@ -0,0 +1,19 @@ +package it.cnr.isti.workflow.manager.executions; + +import java.util.Map; + +public record InteractionResult(boolean completed, Map outputs, Map partialResults) { + + public InteractionResult { + outputs = outputs == null ? Map.of() : Map.copyOf(outputs); + partialResults = partialResults == null ? Map.of() : Map.copyOf(partialResults); + } + + public static InteractionResult completed(Map outputs) { + return new InteractionResult(true, outputs, Map.of()); + } + + public static InteractionResult partial(Map partialResults) { + return new InteractionResult(false, Map.of(), partialResults); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java index 81697a5..0be0f09 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java @@ -5,6 +5,7 @@ import java.util.Map; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.executions.InteractionResult; import it.cnr.isti.workflow.manager.executions.steps.Input; import it.cnr.isti.workflow.manager.flows.model.FlowNode; @@ -32,4 +33,12 @@ public final class NodeExecutors { } throw new IllegalStateException("No executor found for node type " + node.getClass().getName()); } + + public static InteractionResult interact(FlowNode node, List inputs, Map interaction, + Map partialResults, Map authorizations) { + if (node instanceof Block block) { + return BlockExecutors.get(block.getType()).interact((Block) block, inputs, interaction, partialResults, authorizations); + } + throw new IllegalStateException("No interactive executor found for node type " + node.getClass().getName()); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/BlockExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/BlockExecutor.java index e23b92f..8ff48e4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/BlockExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/BlockExecutor.java @@ -4,6 +4,7 @@ import java.util.List; import java.util.Map; import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.executions.InteractionResult; import it.cnr.isti.workflow.manager.blocks.types.BlockType; import it.cnr.isti.workflow.manager.executions.steps.Input; @@ -22,4 +23,9 @@ public interface BlockExecutor { throw new UnsupportedOperationException("Simulation not supported for this block type"); } + default InteractionResult interact(Block block, List inputs, Map interaction, + Map partialResults, Map authorizations) { + throw new UnsupportedOperationException("Interaction not supported for this block type"); + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java new file mode 100644 index 0000000..f16d55b --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java @@ -0,0 +1,168 @@ +package it.cnr.isti.workflow.manager.executions.executors.blocks; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.factories.ChatInteractionBlockFactory; +import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; +import it.cnr.isti.workflow.manager.executions.InteractionResult; +import it.cnr.isti.workflow.manager.executions.steps.Input; +import it.cnr.isti.workflow.manager.llms.ChatMessage; +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; + +@Component +public class ChatInteractionExecutor implements BlockExecutor { + + @Autowired + private Map llmProviders; + + @Override + public Map execute(Block block, List inputs, + Map authorizations) { + throw new UnsupportedOperationException("ChatInteraction blocks require user interaction."); + } + + @Override + public Class getBlockType() { + return ChatInteractionBlockType.class; + } + + @Override + public boolean isInteractive() { + return true; + } + + @Override + public InteractionResult interact(Block block, List inputs, + Map interaction, Map partialResults, Map authorizations) { + ChatInteractionBlockConfiguration configuration = + (ChatInteractionBlockConfiguration) block.getSpecificConfiguration(); + LLMDescriptor llmDescriptor = configuration.getLlmDescriptor(); + LLMProvider llmProvider = resolveProvider(llmDescriptor.provider()); + String authKey = llmProvider.authorizationKey(); + if (llmProvider.requiresAuthorization() + && (!authorizations.containsKey(authKey) || !StringUtils.hasText(String.valueOf(authorizations.get(authKey))))) { + throw new IllegalArgumentException("Missing authorization for provider: " + llmDescriptor.provider()); + } + + if (interaction.containsKey(ChatInteractionBlockFactory.INTERACTION_FIELD)) { + Object messageValue = interaction.get(ChatInteractionBlockFactory.INTERACTION_FIELD); + if (!(messageValue instanceof String message) || !StringUtils.hasText(message)) { + throw new IllegalArgumentException("Missing interaction value for field: " + + ChatInteractionBlockFactory.INTERACTION_FIELD); + } + + String resolvedMessage = resolvePlaceholders(message, inputs); + List history = existingHistory(partialResults); + List messages = history.stream() + .map(this::parseHistoryLine) + .collect(Collectors.toList()); + messages.add(new ChatMessage(ChatMessage.Role.USER, resolvedMessage)); + + String assistantResponse = llmProvider.requiresAuthorization() + ? llmProvider.chat(llmDescriptor.model(), List.copyOf(messages), String.valueOf(authorizations.get(authKey))) + : llmProvider.chat(llmDescriptor.model(), List.copyOf(messages)); + + List updatedHistory = new ArrayList<>(history); + updatedHistory.add(formatConversationLine(ChatMessage.Role.USER, resolvedMessage)); + updatedHistory.add(formatConversationLine(ChatMessage.Role.ASSISTANT, assistantResponse)); + + return InteractionResult.partial(Map.of( + ChatInteractionBlockFactory.RESPONSE_OUTPUT, assistantResponse, + ChatInteractionBlockFactory.HISTORY_OUTPUT, List.copyOf(updatedHistory))); + } + + if (interaction.containsKey(ChatInteractionBlockFactory.FINAL_RESPONSE_FIELD)) { + Object responseValue = interaction.get(ChatInteractionBlockFactory.FINAL_RESPONSE_FIELD); + if (!(responseValue instanceof String response) || !StringUtils.hasText(response)) { + throw new IllegalArgumentException("Missing interaction value for field: " + + ChatInteractionBlockFactory.FINAL_RESPONSE_FIELD); + } + List history = existingHistory(partialResults); + return InteractionResult.completed(Map.of( + ChatInteractionBlockFactory.RESPONSE_OUTPUT, response, + ChatInteractionBlockFactory.HISTORY_OUTPUT, List.copyOf(history))); + } + + throw new IllegalArgumentException("Unsupported interaction field for ChatInteraction"); + } + + private LLMProvider resolveProvider(String providerName) { + LLMProvider llmProvider = llmProviders.get(providerName); + if (llmProvider == null) { + llmProvider = llmProviders.values().stream() + .filter(candidate -> candidate.getName().equals(providerName)) + .findFirst() + .orElse(null); + } + if (llmProvider == null) { + throw new IllegalArgumentException("Provider not found: " + providerName); + } + return llmProvider; + } + + private String formatConversationLine(ChatMessage.Role role, String content) { + return "[" + role.name() + "] " + content; + } + + @SuppressWarnings("unchecked") + private List existingHistory(Map partialResults) { + Object value = partialResults.get(ChatInteractionBlockFactory.HISTORY_OUTPUT); + if (value instanceof List list) { + return ((List) list).stream().map(String::valueOf).toList(); + } + return List.of(); + } + + private ChatMessage parseHistoryLine(String line) { + if (!StringUtils.hasText(line)) { + return new ChatMessage(ChatMessage.Role.USER, ""); + } + if (line.startsWith("[") && line.contains("]")) { + int endRole = line.indexOf(']'); + String rawRole = line.substring(1, endRole); + String content = line.substring(endRole + 1).stripLeading(); + try { + return new ChatMessage(ChatMessage.Role.valueOf(rawRole), content); + } catch (IllegalArgumentException ignored) { + return new ChatMessage(ChatMessage.Role.USER, line); + } + } + return new ChatMessage(ChatMessage.Role.USER, line); + } + + private String resolvePlaceholders(String template, List inputs) { + if (!StringUtils.hasText(template)) { + return template; + } + Map values = new LinkedHashMap<>(); + for (Input input : inputs) { + values.put(input.getDescriptor().getName(), formatInputValue(input.getValue())); + } + String resolved = template; + for (Map.Entry entry : values.entrySet()) { + resolved = resolved.replace("${{" + entry.getKey() + "}}", entry.getValue()); + } + return resolved; + } + + private String formatInputValue(Object value) { + if (value instanceof Collection collection) { + return collection.stream() + .map(item -> item == null ? "null" : item.toString()) + .collect(Collectors.joining(System.lineSeparator())); + } + return value == null ? "null" : value.toString(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/HumanInteractionExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/HumanInteractionExecutor.java index 177420a..9d53a1f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/HumanInteractionExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/HumanInteractionExecutor.java @@ -11,6 +11,7 @@ import org.springframework.util.StringUtils; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; +import it.cnr.isti.workflow.manager.executions.InteractionResult; import it.cnr.isti.workflow.manager.executions.steps.Input; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @@ -75,4 +76,14 @@ public class HumanInteractionExecutor implements BlockExecutor block, List inputs, + Map interaction, Map partialResults, Map authorizations) { + Object value = interaction.get("output"); + if (value == null) { + throw new IllegalArgumentException("Missing interaction value for field: output"); + } + return InteractionResult.completed(Map.of("output", value)); + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java index 7994b01..8f4fed7 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java @@ -1,6 +1,7 @@ package it.cnr.isti.workflow.manager.executions.steps; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; @@ -38,6 +39,10 @@ public class Step implements InputListener { @Getter private final List outputs = new ArrayList<>(); + @Getter + @JsonIgnore + private final Map partialResults = new HashMap<>(); + @Getter private StepStatus status = StepStatus.WAITING_FOR_INPUT; @@ -150,15 +155,27 @@ public class Step implements InputListener { logger.info("Resuming step " + this.id + " of node " + this.node.getName()); this.status = StepStatus.RUNNING; listener.resumed(this.id); + var interactionResult = NodeExecutors.interact(this.node, this.inputs, providedOutputs, Map.copyOf(this.partialResults), authorizations); + this.partialResults.clear(); + this.partialResults.putAll(interactionResult.partialResults()); + listener.partialUpdated(this.id, this.partialResults); + + if (!interactionResult.completed()) { + this.status = StepStatus.WAITING_FOR_INTERACTION; + listener.paused(this.id); + return; + } + Map resolvedOutputs = interactionResult.outputs(); for (Output output : this.outputs) { - if (providedOutputs.containsKey(output.getDescriptor().getName())) { - output.setValue(providedOutputs.get(output.getDescriptor().getName())); + if (resolvedOutputs.containsKey(output.getDescriptor().getName())) { + output.setValue(resolvedOutputs.get(output.getDescriptor().getName())); } else { throw new IllegalArgumentException("Missing output value for: " + output.getDescriptor().getName()); } } + this.partialResults.clear(); this.status = StepStatus.COMPLETED; - listener.completed(this.id, providedOutputs); + listener.completed(this.id, resolvedOutputs); } private void refreshStateFromInputs() { 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 487b49f..13601be 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 @@ -115,26 +115,23 @@ public class FlowDataValidator implements ConstraintValidator containerConfiguration = container.getSpecificConfiguration(); FlowData subFlow = containerConfiguration.getSubFlow(); - if (subFlow == null || subFlow.getNodes().isEmpty()) { - throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Container subFlow must contain at least one node")); - } + if (subFlow != null && !subFlow.getNodes().isEmpty()) { + List nestedNodes = subFlow.getNodes(); + if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container)) { + throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", + "Nested containers are not supported")); + } + if (nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) { + throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", + "Interactive blocks inside containers are not supported yet")); + } - List nestedNodes = subFlow.getNodes(); - if (nestedNodes.stream().anyMatch(candidate -> candidate instanceof Container)) { - throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Nested containers are not supported")); - } - if (nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) { - throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Interactive blocks inside containers are not supported yet")); - } - - try { - validateFlowData(subFlow); - } catch (FlowValidationException e) { - throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", - "Invalid container subFlow: " + e.getMessage())); + try { + validateFlowData(subFlow); + } catch (FlowValidationException e) { + throw validationError(error("container", container.getId(), "specificConfiguration.subFlow", + "Invalid container subFlow: " + e.getMessage())); + } } final Container canonicalContainer; 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 f7e5738..77e0552 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 @@ -67,6 +67,33 @@ public class FlowExecutionValidator { if (container.getSpecificConfiguration() instanceof ContainerConfiguration containerConfiguration && containerConfiguration.getSubFlow() != null) { + FlowData subFlow = containerConfiguration.getSubFlow(); + if (subFlow.getNodes().isEmpty()) { + errors.add(new ValidationError( + "container", + container.getId(), + "specificConfiguration.subFlow", + "Container subFlow must contain at least one node")); + } else { + subFlow.getNodes().stream() + .filter(candidate -> candidate instanceof Container) + .findFirst() + .ifPresent(candidate -> errors.add(new ValidationError( + "container", + container.getId(), + "specificConfiguration.subFlow", + "Nested containers are not supported"))); + + subFlow.getNodes().stream() + .filter(FlowNode::isUserInteractive) + .findFirst() + .ifPresent(candidate -> errors.add(new ValidationError( + "container", + container.getId(), + "specificConfiguration.subFlow", + "Interactive nodes inside containers are not supported yet"))); + } + errors.addAll(collectErrors(containerConfiguration.getSubFlow()).stream() .map(error -> new ValidationError( "container", diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/ChatMessage.java b/src/main/java/it/cnr/isti/workflow/manager/llms/ChatMessage.java new file mode 100644 index 0000000..a56aef1 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/ChatMessage.java @@ -0,0 +1,17 @@ +package it.cnr.isti.workflow.manager.llms; + +import java.util.Objects; + +public record ChatMessage(Role role, String content) { + + public ChatMessage { + role = role == null ? Role.USER : role; + content = Objects.requireNonNullElse(content, ""); + } + + public enum Role { + SYSTEM, + USER, + ASSISTANT + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/LLMProvider.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/LLMProvider.java index 3def45e..c986e35 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/LLMProvider.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/LLMProvider.java @@ -1,6 +1,9 @@ package it.cnr.isti.workflow.manager.llms.providers; import java.util.List; +import java.util.stream.Collectors; + +import it.cnr.isti.workflow.manager.llms.ChatMessage; public interface LLMProvider { @@ -12,6 +15,16 @@ public interface LLMProvider { return generate(model, prompt); } + default String chat(String model, List messages) { + return generate(model, flattenConversation(messages)); + } + + default String chat(String model, List messages, String authorization) { + return requiresAuthorization() + ? generate(model, flattenConversation(messages), authorization) + : chat(model, messages); + } + default boolean requiresAuthorization() { return false; } @@ -28,4 +41,10 @@ public interface LLMProvider { return "Provider authorization value (API key or token)."; } + private static String flattenConversation(List messages) { + return messages == null ? "" : messages.stream() + .map(message -> "[" + message.role().name() + "] " + message.content()) + .collect(Collectors.joining(System.lineSeparator() + System.lineSeparator())); + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/google/GeminiLLMProvider.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/google/GeminiLLMProvider.java index 80d6d5e..691f104 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/google/GeminiLLMProvider.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/google/GeminiLLMProvider.java @@ -14,6 +14,8 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import it.cnr.isti.workflow.manager.llms.ChatMessage; + @Service public class GeminiLLMProvider implements LLMProvider { @@ -40,23 +42,22 @@ public class GeminiLLMProvider implements LLMProvider { } @Override - public String generate(String model, String prompt, String authorization) { - Objects.requireNonNull(prompt, "prompt cannot be null"); - Objects.requireNonNull(model, "model cannot be null"); + public String chat(String model, List messages, String authorization) { + Objects.requireNonNull(messages, "messages cannot be null"); Objects.requireNonNull(authorization, "authorization cannot be null"); if (!MODELS.contains(model)) throw new IllegalArgumentException("Model not supported: " + model); Map requestBody = Map.of( - "contents", List.of( // Usa List.of invece di new Object[] - Map.of( - "role", "user", - "parts", List.of(Map.of("text", prompt)) // Usa List.of per "parts" - ))); + "contents", messages.stream() + .map(message -> Map.of( + "role", toGeminiRole(message.role()), + "parts", List.of(Map.of("text", message.content())))) + .toList()); WebClient webClient = webClientBuilder.baseUrl(GEMINI_URL).build(); Mono result = webClient.post() - .uri(uriBuilder -> uriBuilder.path("v1beta/models/gemini-2.0-flash:generateContent") + .uri(uriBuilder -> uriBuilder.path("v1beta/models/" + model + ":generateContent") .queryParam("key", authorization) .build()) .contentType(MediaType.APPLICATION_JSON) @@ -79,6 +80,11 @@ public class GeminiLLMProvider implements LLMProvider { return result.block(); } + @Override + public String generate(String model, String prompt, String authorization) { + return chat(model, List.of(new ChatMessage(ChatMessage.Role.USER, prompt)), authorization); + } + @Override public boolean requiresAuthorization() { return true; @@ -94,4 +100,11 @@ public class GeminiLLMProvider implements LLMProvider { return MODELS; } + private String toGeminiRole(ChatMessage.Role role) { + return switch (role) { + case SYSTEM, USER -> "user"; + case ASSISTANT -> "model"; + }; + } + } diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java index 36c0ae9..aed1151 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/InternalOllamaLLMProvider.java @@ -12,7 +12,9 @@ import org.springframework.stereotype.Service; import org.springframework.web.reactive.function.client.WebClient; import com.fasterxml.jackson.databind.ObjectMapper; +import it.cnr.isti.workflow.manager.llms.ChatMessage; import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; +import it.cnr.isti.workflow.manager.llms.providers.ollama.response.ChatResponse; import it.cnr.isti.workflow.manager.llms.providers.ollama.response.GenerateResponse; import it.cnr.isti.workflow.manager.llms.providers.ollama.response.ModelInfo; import it.cnr.isti.workflow.manager.llms.providers.ollama.response.ModelResponse; @@ -74,6 +76,40 @@ public class InternalOllamaLLMProvider implements LLMProvider { return parseGenerateResponse(result.block(), mapper); } + @Override + public String chat(String model, List messages) { + Objects.requireNonNull(messages, "messages cannot be null"); + Objects.requireNonNull(model, "model cannot be null"); + ObjectMapper mapper = new ObjectMapper(); + Map bodyMap = Map.of( + "model", model, + "messages", messages.stream() + .map(message -> Map.of( + "role", message.role().name().toLowerCase(), + "content", message.content())) + .toList(), + "stream", false); + + WebClient webClient = webClientBuilder.baseUrl(this.ollamaURL).build(); + + Mono result = webClient.post() + .uri(uriBuilder -> uriBuilder.pathSegment("chat").build()) + .header("Authorization", "Bearer " + ollamaKey) + .contentType(MediaType.APPLICATION_JSON) + .bodyValue(bodyMap) + .retrieve() + .onStatus( + status -> status.is5xxServerError(), + clientResponse -> clientResponse.bodyToMono(String.class) + .defaultIfEmpty("Error: server without body") + .flatMap(body -> Mono.error(new RuntimeException( + "Error 5xx: " + body)))) + .bodyToMono(String.class) + .timeout(Duration.ofMinutes(2)); + + return parseChatResponse(result.block(), mapper); + } + private String parseGenerateResponse(String responseBody, ObjectMapper mapper) { if (responseBody == null || responseBody.isBlank()) { throw new RuntimeException("Empty response body from Ollama generate endpoint"); @@ -96,6 +132,28 @@ public class InternalOllamaLLMProvider implements LLMProvider { return trimmedResponse; } + private String parseChatResponse(String responseBody, ObjectMapper mapper) { + if (responseBody == null || responseBody.isBlank()) { + throw new RuntimeException("Empty response body from Ollama chat endpoint"); + } + + String trimmedResponse = responseBody.trim(); + if (trimmedResponse.startsWith("{")) { + try { + ChatResponse response = mapper.readValue(trimmedResponse, ChatResponse.class); + if (response.getMessage() == null || response.getMessage().getContent() == null) { + throw new RuntimeException("Missing 'message.content' field in Ollama chat response"); + } + return response.getMessage().getContent(); + } catch (Exception e) { + throw new RuntimeException("Unable to parse Ollama chat response", e); + } + } + + log.debug("Ollama chat endpoint returned text/plain response"); + return trimmedResponse; + } + public List getRegisteredModels() { WebClient webClient = webClientBuilder.baseUrl(this.ollamaURL).build(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponse.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponse.java new file mode 100644 index 0000000..db869c5 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponse.java @@ -0,0 +1,13 @@ +package it.cnr.isti.workflow.manager.llms.providers.ollama.response; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; + +import lombok.Data; +import lombok.NoArgsConstructor; + +@JsonIgnoreProperties(ignoreUnknown = true) +@Data +@NoArgsConstructor +public class ChatResponse { + private ChatResponseMessage message; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponseMessage.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponseMessage.java new file mode 100644 index 0000000..3e07403 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/ollama/response/ChatResponseMessage.java @@ -0,0 +1,14 @@ +package it.cnr.isti.workflow.manager.llms.providers.ollama.response; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; + +import lombok.Data; +import lombok.NoArgsConstructor; + +@JsonIgnoreProperties(ignoreUnknown = true) +@Data +@NoArgsConstructor +public class ChatResponseMessage { + private String role; + private String content; +} diff --git a/src/main/resources/workflow-editor-init/flows.json b/src/main/resources/workflow-editor-init/flows.json index 6a1fdf7..341fd94 100644 --- a/src/main/resources/workflow-editor-init/flows.json +++ b/src/main/resources/workflow-editor-init/flows.json @@ -353,5 +353,320 @@ } ] } + }, + { + "name": "Resume Analysis Pipeline", + "description": "Analyze a curriculum vitae through multiple LLM steps and produce a final recruiter-oriented assessment without human interaction.", + "owner": "testuser", + "createdAt": "2026-03-18T09:30:00", + "lastUpdateAt": "2026-03-18T09:30:00", + "published": false, + "finalized": false, + "flow": { + "blocks": [ + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da1001", + "position": { + "x": -1200, + "y": 0 + }, + "name": "normalize-resume", + "inputs": [ + { + "name": "resume", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + }, + { + "type": "TEXT", + "multiple": true + } + ] + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "normalize-resume", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "You are a recruiting analyst. Normalize the following curriculum vitae into a clean structured profile with sections for identity, role, years of experience, skills, experiences, education, certifications, and languages. Resume: ${{resume}}" + }, + "typeName": "LLMBlock" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da1002", + "position": { + "x": -600, + "y": -220 + }, + "name": "extract-core-skills", + "inputs": [ + { + "name": "profile", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + }, + { + "type": "TEXT", + "multiple": true + } + ] + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "extract-core-skills", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "From this normalized candidate profile, extract the 10 most relevant professional skills for recruiting evaluation. For each skill include strength level, evidence, and business relevance. Profile: ${{profile}}" + }, + "typeName": "LLMBlock" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da1003", + "position": { + "x": -600, + "y": 0 + }, + "name": "summarize-experience", + "inputs": [ + { + "name": "profile", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + }, + { + "type": "TEXT", + "multiple": true + } + ] + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "summarize-experience", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "Summarize the candidate's professional experience from this normalized profile. Highlight seniority, domains, progression, leadership, and most relevant achievements. Profile: ${{profile}}" + }, + "typeName": "LLMBlock" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da1004", + "position": { + "x": -600, + "y": 220 + }, + "name": "detect-risks-and-gaps", + "inputs": [ + { + "name": "profile", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + }, + { + "type": "TEXT", + "multiple": true + } + ] + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "detect-risks-and-gaps", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "Review this normalized candidate profile and identify hiring risks, possible gaps, unclear claims, missing evidence, and follow-up questions a recruiter should ask. Profile: ${{profile}}" + }, + "typeName": "LLMBlock" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da1005", + "position": { + "x": 100, + "y": 0 + }, + "name": "build-final-assessment", + "inputs": [ + { + "name": "skillsAnalysis", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + }, + { + "name": "experienceSummary", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + }, + { + "name": "riskAnalysis", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "outputs": [ + { + "name": "response", + "type": "TEXT", + "multiple": false, + "valueKinds": [ + { + "type": "TEXT", + "multiple": false + } + ] + } + ], + "specificConfiguration": { + "type": "LLMBlockConfiguration", + "name": "build-final-assessment", + "llmDescriptor": { + "provider": "InternalOllama", + "model": "gemma:7b" + }, + "prompt": "Create a recruiter-facing final assessment using these three analyses. Return: executive summary, strengths, concerns, seniority estimate, best-fit roles, and recommendation (strong yes / yes / maybe / no). Skills: ${{skillsAnalysis}} Experience: ${{experienceSummary}} Risks: ${{riskAnalysis}}" + }, + "typeName": "LLMBlock" + } + ], + "connections": [ + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2001", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1001", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1002", + "targetName": "profile" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2002", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1001", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1003", + "targetName": "profile" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2003", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1001", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1004", + "targetName": "profile" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2004", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1002", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1005", + "targetName": "skillsAnalysis" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2005", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1003", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1005", + "targetName": "experienceSummary" + }, + { + "id": "cc1f2b22-1b38-4a53-a4f3-4bde50da2006", + "sourceId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1004", + "sourceName": "response", + "targetId": "cc1f2b22-1b38-4a53-a4f3-4bde50da1005", + "targetName": "riskAnalysis" + } + ] + } } ] diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java index cbdca4b..ecbc9b7 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/BlocksControllerTest.java @@ -5,6 +5,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertThrows; import java.util.List; @@ -17,7 +18,10 @@ 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.configurations.ChatInteractionBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionInput; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; +import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; 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; @@ -79,6 +83,132 @@ public class BlocksControllerTest { assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("response"))); } + @Test + public void createChatInteractionBlockExposesInteractiveOutputs() { + ChatInteractionBlockConfiguration config = ChatInteractionBlockConfiguration.builder() + .name("Recruiter chat") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .inputs(List.of(new ChatInteractionInput("candidate", it.cnr.isti.workflow.manager.ios.IOType.TEXT))) + .build(); + + Block block = blocksController.create(config); + + assertNotNull(block); + assertEquals(ChatInteractionBlockType.TYPE, block.getType().getName()); + assertTrue(block.getInputs().stream().anyMatch(input -> input.getName().equals("candidate"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("response") && !output.isMultiple())); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("history") && output.isMultiple())); + } + + @Test + public void createChatInteractionBlockAllowsPartialConfigurationForUiDraft() { + Block block = blocksController.create(ChatInteractionBlockConfiguration.empty()); + + assertNotNull(block); + assertEquals(ChatInteractionBlockType.TYPE, block.getType().getName()); + assertEquals(ChatInteractionBlockType.TYPE, block.getName()); + assertTrue(block.getInputs().isEmpty()); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("response"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("history"))); + } + + @Test + public void createChatInteractionBlockDefaultsMissingInputTypeToText() { + ChatInteractionBlockConfiguration config = ChatInteractionBlockConfiguration.builder() + .name("Recruiter chat") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .inputs(List.of(new ChatInteractionInput("candidate", null))) + .build(); + + Block block = blocksController.create(config); + + assertNotNull(block); + assertEquals(it.cnr.isti.workflow.manager.ios.IOType.TEXT, + block.getInputs().stream() + .filter(input -> input.getName().equals("candidate")) + .findFirst() + .orElseThrow() + .getType()); + } + + @Test + public void createChatInteractionBlockAcceptsLegacyInputTypeAlias() throws Exception { + String payload = """ + { + "type": "ChatInteractionBlockConfiguration", + "name": "Recruiter chat", + "llmDescriptor": { + "provider": "testProvider", + "model": "testModel" + }, + "inputs": [ + { + "name": "candidate", + "type": "TEXT", + "multiple": true + } + ] + } + """; + + ChatInteractionBlockConfiguration config = ObjectMapperHolder.mapper.readValue( + payload, + ChatInteractionBlockConfiguration.class); + Block block = blocksController.create(config); + + assertEquals(it.cnr.isti.workflow.manager.ios.IOType.TEXT, + block.getInputs().stream() + .filter(input -> input.getName().equals("candidate")) + .findFirst() + .orElseThrow() + .getType()); + assertTrue(block.getInputs().stream() + .filter(input -> input.getName().equals("candidate")) + .findFirst() + .orElseThrow() + .isMultiple()); + } + + @Test + public void createChatInteractionBlockRejectsDuplicateInputNames() { + ChatInteractionBlockConfiguration config = ChatInteractionBlockConfiguration.builder() + .name("Recruiter chat") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .inputs(List.of( + new ChatInteractionInput("candidate", it.cnr.isti.workflow.manager.ios.IOType.TEXT), + new ChatInteractionInput("candidate", it.cnr.isti.workflow.manager.ios.IOType.TEXT, true))) + .build(); + + IllegalArgumentException exception = assertThrows(IllegalArgumentException.class, + () -> blocksController.create(config)); + assertTrue(exception.getMessage().contains("inputs must have unique names")); + } + + @Test + public void createChatInteractionBlockRejectsNonTextInputType() { + ChatInteractionBlockConfiguration config = ChatInteractionBlockConfiguration.builder() + .name("Recruiter chat") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .inputs(List.of(new ChatInteractionInput("candidate", it.cnr.isti.workflow.manager.ios.IOType.FILE))) + .build(); + + IllegalArgumentException exception = assertThrows(IllegalArgumentException.class, + () -> blocksController.create(config)); + assertTrue(exception.getMessage().contains("ChatInteraction inputs support only TEXT or TEXT[]")); + } + @Test public void LLMSchemaHasLongTextHints() { BlockConfigurationDescriptor descriptor = blocksController @@ -102,6 +232,68 @@ public class BlocksControllerTest { assertFalse(prompt.has("x-ui-rows")); } + @Test + public void chatInteractionSchemaDeclaresUniqueInputNames() { + BlockConfigurationDescriptor descriptor = blocksController + .getConfigurationDescriptorForType(ChatInteractionBlockType.TYPE); + + JsonNode schema = (JsonNode) descriptor.schema(); + JsonNode inputs = schema.path("properties").path("inputs"); + assertEquals("name", inputs.path("x-ui-unique-by").asText()); + JsonNode definitions = schema.has("definitions") ? schema.path("definitions") : schema.path("$defs"); + JsonNode inputDefinition = definitions.fields().next().getValue(); + for (java.util.Iterator> it = definitions.fields(); it.hasNext();) { + java.util.Map.Entry entry = it.next(); + if (entry.getKey().contains("ChatInteractionInput")) { + inputDefinition = entry.getValue(); + break; + } + } + JsonNode itemProperties = inputDefinition.path("properties"); + assertEquals(8, itemProperties.path("name").path("maxLength").asInt()); + assertTrue(itemProperties.has("ioType")); + assertTrue(itemProperties.has("multiple")); + assertFalse(itemProperties.has("type")); + assertEquals("TEXT", itemProperties.path("ioType").path("enum").get(0).asText()); + assertEquals(1, itemProperties.path("ioType").path("enum").size()); + assertEquals("type", itemProperties.path("ioType").path("x-ui-label").asText()); + assertFalse(itemProperties.path("ioType").has("x-ui-description")); + } + + @Test + public void chatInteractionDescriptorExposesInteractionContract() { + BlockConfigurationDescriptor descriptor = blocksController + .getConfigurationDescriptorForType(ChatInteractionBlockType.TYPE); + + assertNotNull(descriptor.interactionContract()); + assertEquals("chat-session", descriptor.interactionContract().kind()); + assertEquals("message", descriptor.interactionContract().messageField()); + assertEquals("response", descriptor.interactionContract().completionField()); + assertEquals("history", descriptor.interactionContract().historyField()); + assertEquals("response", descriptor.interactionContract().responseField()); + assertTrue(descriptor.interactionContract().supportsPartialResult()); + } + + @Test + public void humanInteractionDescriptorExposesInteractionContract() { + BlockConfigurationDescriptor descriptor = blocksController + .getConfigurationDescriptorForType(HumanInteractionBlockType.TYPE); + + assertNotNull(descriptor.interactionContract()); + assertEquals("single-response", descriptor.interactionContract().kind()); + assertEquals("output", descriptor.interactionContract().completionField()); + assertEquals("output", descriptor.interactionContract().responseField()); + assertFalse(descriptor.interactionContract().supportsPartialResult()); + } + + @Test + public void llmDescriptorDoesNotExposeInteractionContract() { + BlockConfigurationDescriptor descriptor = blocksController + .getConfigurationDescriptorForType(LLMBlockType.TYPE); + + assertNull(descriptor.interactionContract()); + } + @Test public void getLLMExampleForType() { Block block = blocksController.getExampleForType(LLMBlockType.TYPE); @@ -129,6 +321,28 @@ public class BlocksControllerTest { .anyMatch(capability -> capability.type() == IOCapabilityType.TEXT && !capability.multiple())); } + @Test + public void getChatInteractionExampleForType() { + Block block = blocksController.getExampleForType(ChatInteractionBlockType.TYPE); + + assertNotNull(block); + assertEquals(ChatInteractionBlockType.TYPE, block.getType().getName()); + assertEquals(ChatInteractionBlockType.TYPE, block.getName()); + assertNotNull(block.getSpecificConfiguration()); + assertTrue(block.getInputs().isEmpty()); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("response"))); + assertTrue(block.getOutputs().stream().anyMatch(output -> output.getName().equals("history") && output.isMultiple())); + } + + @Test + public void getLegacyChatHumanInteractionExampleForType() { + Block block = blocksController.getExampleForType("ChatHumanInteraction"); + + assertNotNull(block); + assertEquals(ChatInteractionBlockType.TYPE, block.getType().getName()); + assertEquals(ChatInteractionBlockType.TYPE, block.getName()); + } + @Test public void createLlmBlockWithoutConfiguredPromptUsesPromptInput() { Block block = blocksController.create(LLMBlockConfiguration.builder() 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 812dbd5..dfdae84 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 @@ -3,7 +3,6 @@ package it.cnr.isti.workflow.manager.controllers; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; -import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import java.util.List; @@ -12,7 +11,6 @@ import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.TestPropertySource; -import org.springframework.web.server.ResponseStatusException; import com.fasterxml.jackson.databind.JsonNode; @@ -100,6 +98,22 @@ public class ContainersControllerTest { assertTrue(subFlow.path("x-ui-structural").asBoolean()); } + @Test + public void iteratorContainerSchemaMarksIterationInputAsEnabledWhenSubFlowIsPresent() { + ContainersController.ContainerConfigurationDescriptor descriptor = containersController.getTypes().stream() + .filter(type -> IteratorContainerType.TYPE.equals(type.type())) + .findFirst() + .orElseThrow(); + + JsonNode schema = (JsonNode) descriptor.schema(); + JsonNode iterationInput = schema.path("properties").path("iterationInput"); + assertEquals("subFlow", iterationInput.path("x-ui-enabled-when").path("field").asText()); + assertTrue(iterationInput.path("x-ui-enabled-when").path("present").asBoolean()); + assertEquals("inputs", iterationInput.path("x-ui-options-from-node").path("collection").asText()); + assertEquals("name", iterationInput.path("x-ui-options-from-node").path("valueField").asText()); + assertEquals("name", iterationInput.path("x-ui-options-from-node").path("labelField").asText()); + } + @Test public void createGenericContainerExposesOpenHandles() { Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() @@ -227,15 +241,14 @@ public class ContainersControllerTest { } @Test - public void createGenericContainerRejectsEmptySubFlow() { - ResponseStatusException exception = assertThrows(ResponseStatusException.class, - () -> containersController.create(GenericContainerConfiguration.builder() - .name("Container") - .build())); + public void createGenericContainerAllowsEmptySubFlow() { + Container container = containersController.create(GenericContainerConfiguration.builder() + .name("Container") + .build()); - assertEquals(400, exception.getStatusCode().value()); - assertTrue(exception.getReason().contains("Invalid parameter subFlow")); - assertTrue(exception.getReason().contains("must contain at least one node")); + assertNotNull(container); + assertTrue(container.getInputs().isEmpty()); + assertTrue(container.getOutputs().isEmpty()); } @Test @@ -256,12 +269,12 @@ public class ContainersControllerTest { .build()); assertNotNull(container); - assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidates") && input.isMultiple())); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate") && input.isMultiple())); assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response") && output.isMultiple())); } @Test - public void createIteratorContainerRejectsUnknownIterationInput() { + public void createIteratorContainerAllowsUnknownIterationInputDuringCreate() { Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() .name("Analyze") .llmDescriptor(LLMDescriptor.builder() @@ -271,14 +284,59 @@ public class ContainersControllerTest { .prompt("Analyze ${{candidate}}") .build()); - ResponseStatusException exception = assertThrows(ResponseStatusException.class, - () -> containersController.create(IteratorContainerConfiguration.builder() - .name("Iterator") - .subFlow(FlowData.builder().block(internalBlock).build()) - .iterationInput("unknown") - .build())); + Container container = 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")); + assertNotNull(container); + assertEquals("candidate", ((IteratorContainerConfiguration) container.getSpecificConfiguration()).getIterationInput()); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate") && input.isMultiple())); + assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response") && output.isMultiple())); + } + + @Test + public void createIteratorContainerWithoutIterationInputDefaultsIterationInputWhenOnlyOneChoiceExists() { + 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()) + .build()); + + assertNotNull(container); + assertEquals("candidate", ((IteratorContainerConfiguration) container.getSpecificConfiguration()).getIterationInput()); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("candidate") && input.isMultiple())); + assertTrue(container.getOutputs().stream().anyMatch(output -> output.getName().equals("response") && output.isMultiple())); + } + + @Test + public void createIteratorContainerOverridesIterationInputWhenNewSubFlowHasSingleOpenInput() { + Block internalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("Analyze") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .prompt("Analyze ${{person}}") + .build()); + + Container container = containersController.create(IteratorContainerConfiguration.builder() + .name("Iterator") + .subFlow(FlowData.builder().block(internalBlock).build()) + .iterationInput("candidate") + .build()); + + assertNotNull(container); + assertEquals("person", ((IteratorContainerConfiguration) container.getSpecificConfiguration()).getIterationInput()); + assertTrue(container.getInputs().stream().anyMatch(input -> input.getName().equals("person") && input.isMultiple())); } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java index c0f565c..30b68ce 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/ExecutionControllerTest.java @@ -106,7 +106,8 @@ public class ExecutionControllerTest { try { Thread.sleep(1000); } catch (InterruptedException e) { - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } @@ -116,8 +117,7 @@ public class ExecutionControllerTest { try { logger.info("Execution object: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(executionObject)); } catch (JsonProcessingException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + throw new RuntimeException(e); } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java index 6d24186..e57f57b 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/FlowControllerTest.java @@ -31,6 +31,9 @@ import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.factories.ConditionalBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.ConditionalBlockType; 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.IteratorContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; import it.cnr.isti.workflow.manager.flows.FlowTestCreator; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.Flow; @@ -56,6 +59,9 @@ public class FlowControllerTest { @Autowired private BlocksController blocksController; + @Autowired + private ContainersController containersController; + @Autowired private MockMvc mockMvc; @@ -106,7 +112,7 @@ public class FlowControllerTest { try { System.out.println(ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(retrieved)); } catch (JsonProcessingException e) { - e.printStackTrace(); + throw new RuntimeException(e); } assertNotNull(retrieved); @@ -498,6 +504,73 @@ public class FlowControllerTest { assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status()); } + @Test + public void createFlowWithDraftIteratorContainerReturnsDraftStatus() { + Block firstInternalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("First") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .prompt("Analyze ${{candidate}}") + .build()); + Block secondInternalBlock = blocksController.create(LLMBlockConfiguration.builder() + .name("Second") + .llmDescriptor(LLMDescriptor.builder() + .provider("testProvider") + .model("testModel") + .build()) + .prompt("Analyze ${{profile}}") + .build()); + + Container iteratorContainer = containersController.create( + IteratorContainerConfiguration.builder() + .name("Iterator") + .subFlow(FlowData.builder() + .block(firstInternalBlock) + .block(secondInternalBlock) + .build()) + .build()); + + FlowCreateRequest request = new FlowCreateRequest( + "Draft Iterator Flow", + "Container missing iterationInput should still be savable", + FlowData.builder() + .container(iteratorContainer) + .build()); + + ResponseEntity createResponse = flowController.createFlow( + request, + new LoginEntity("testuser", "testpassword")); + + assertTrue(createResponse.getStatusCode().is2xxSuccessful()); + assertNotNull(createResponse.getBody()); + assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status()); + } + + @Test + public void createFlowWithEmptyIteratorContainerReturnsDraftStatus() { + Container iteratorContainer = containersController.create( + IteratorContainerConfiguration.builder() + .name("Iterator") + .build()); + + FlowCreateRequest request = new FlowCreateRequest( + "Empty Iterator Flow", + "Container with empty subFlow should still be savable as draft", + FlowData.builder() + .container(iteratorContainer) + .build()); + + ResponseEntity createResponse = flowController.createFlow( + request, + new LoginEntity("testuser", "testpassword")); + + assertTrue(createResponse.getStatusCode().is2xxSuccessful()); + assertNotNull(createResponse.getBody()); + assertEquals(FlowViewStatus.DRAFT, createResponse.getBody().status()); + } + @Test public void createFlowRejectsConditionalBranchMerge() { LLMDescriptor llmDescriptor = LLMDescriptor.builder() 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 a6c283c..1319656 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 @@ -15,12 +15,15 @@ import org.springframework.boot.test.context.TestConfiguration; import org.springframework.context.annotation.Bean; import org.springframework.test.context.TestPropertySource; import org.springframework.util.ResourceUtils; +import org.springframework.web.server.ResponseStatusException; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; 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.ChatInteractionBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionInput; 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; @@ -39,12 +42,15 @@ import it.cnr.isti.workflow.manager.flows.model.Flow; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.factories.ChatInteractionBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.HTTPServerCallBlockFactory; import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory; +import it.cnr.isti.workflow.manager.blocks.types.ChatInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.ios.IOType; +import it.cnr.isti.workflow.manager.llms.ChatMessage; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @@ -73,6 +79,11 @@ public class ExecutionTest { return "Hello, " + prompt.replace("Hello, ", "").replace("!", "") + "!"; } + + @Override + public String chat(String model, List messages) { + return "Chat, " + model + "!"; + } }; } } @@ -86,6 +97,9 @@ public class ExecutionTest { @Autowired LLMBlockFactory llmBlockFactory; + @Autowired + ChatInteractionBlockFactory chatInteractionBlockFactory; + @Autowired LLMBlockType llmBlockType; @@ -127,14 +141,88 @@ public class ExecutionTest { try { Thread.sleep(100); } catch (InterruptedException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertTrue(eo.getContext().getStatus().isFinalState()); eo.getContext().getResult().forEach((k,v) -> System.out.println("Result: " + k + " -> " + v)); } + @Test + public void createChatInteractionExecutionSetInteractionAndStart() { + Block chatBlock = chatInteractionBlockFactory.create(ChatInteractionBlockConfiguration.builder() + .name("Recruiter Chat") + .llmDescriptor(llmBrick) + .inputs(List.of(new ChatInteractionInput("cand", IOType.TEXT, false))) + .build()); + + Flow flow = Flow.builder() + .name("Chat flow") + .description("Single chat block") + .block(chatBlock) + .build(); + + ExecutionObject execObject = executionsService.createExecution(flow); + execObject = executionsService.prepareInput(execObject.getId(), chatBlock.getId(), "cand", "John Doe"); + + execObject = executionsService.startExecution(execObject.getId()); + while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { + execObject = executionsService.getExecution(execObject.getId()); + } + assertEquals(ExecutionStatus.WAITING, execObject.getContext().getStatus()); + execObject.setInteractionValue(chatBlock.getId(), ChatInteractionBlockFactory.INTERACTION_FIELD, "Hello ${{cand}}"); + execObject = executionsService.getExecution(execObject.getId()); + + assertEquals(ExecutionStatus.WAITING, execObject.getContext().getStatus()); + Object partialResponse = execObject.getContext().getPartialResult() + .get(new FieldKey(chatBlock.getId(), ChatInteractionBlockFactory.RESPONSE_OUTPUT)); + assertEquals("Chat, testModel!", partialResponse); + + Object partialConversation = execObject.getContext().getPartialResult() + .get(new FieldKey(chatBlock.getId(), ChatInteractionBlockFactory.HISTORY_OUTPUT)); + assertTrue(partialConversation instanceof List); + assertEquals(2, ((List) partialConversation).size()); + assertTrue(((List) partialConversation).contains("[USER] Hello John Doe")); + assertTrue(((List) partialConversation).contains("[ASSISTANT] Chat, testModel!")); + assertTrue(execObject.getContext().getResult().isEmpty()); + + execObject.setInteractionValue(chatBlock.getId(), ChatInteractionBlockFactory.FINAL_RESPONSE_FIELD, + "Candidate approved"); + execObject = executionsService.getExecution(execObject.getId()); + + assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); + Object response = execObject.getContext().getResult() + .get(new FieldKey(chatBlock.getId(), ChatInteractionBlockFactory.RESPONSE_OUTPUT)); + assertEquals("Candidate approved", response); + + Object conversation = execObject.getContext().getResult() + .get(new FieldKey(chatBlock.getId(), ChatInteractionBlockFactory.HISTORY_OUTPUT)); + assertTrue(conversation instanceof List); + assertEquals(2, ((List) conversation).size()); + assertTrue(((List) conversation).contains("[USER] Hello John Doe")); + assertTrue(((List) conversation).contains("[ASSISTANT] Chat, testModel!")); + assertTrue(execObject.getContext().getPartialResult().isEmpty()); + } + + @Test + public void createExecutionRejectsChatInteractionWithoutLlmDescriptor() { + Block chatBlock = chatInteractionBlockFactory + .create(ChatInteractionBlockConfiguration.empty()); + + Flow flow = Flow.builder() + .name("Invalid Chat Flow") + .description("Chat flow without llmDescriptor") + .block(chatBlock) + .build(); + + ResponseStatusException exception = org.junit.jupiter.api.Assertions.assertThrows( + ResponseStatusException.class, + () -> executionsService.createExecution(flow)); + + assertTrue(exception.getReason().contains("\"field\":\"specificConfiguration.llmDescriptor\"")); + } + @Test public void createInteractiveExecutionSetInputAndStart() { ExecutionObject eo = createInteractiveExecutionAndSetInputInternally(); @@ -145,8 +233,8 @@ public class ExecutionTest { try { Thread.sleep(100); } catch (InterruptedException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertFalse(eo.getContext().getStatus().isFinalState()); @@ -353,7 +441,7 @@ public class ExecutionTest { 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")); + executionsService.prepareInput(execObject.getId(), container.getId(), "name", List.of("Alice", "Bob")); execObject = executionsService.startExecution(execObject.getId()); while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { try { diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionWithContainer.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionWithContainer.java index bf413f0..28ad5fc 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionWithContainer.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionWithContainer.java @@ -80,7 +80,8 @@ public class ExecutionWithContainer { try { Thread.sleep(100); } catch (InterruptedException e) { - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertTrue(execObject.getContext().getStatus().isFinalState()); @@ -113,7 +114,8 @@ public class ExecutionWithContainer { try { Thread.sleep(100); } catch (InterruptedException e) { - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertTrue(execObject.getContext().getStatus().isFinalState()); @@ -145,7 +147,8 @@ public class ExecutionWithContainer { Thread.sleep(100); System.out.println("Execution status: " + execObject.getContext().getStatus()); } catch (InterruptedException e) { - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertTrue(execObject.getContext().getStatus().isFinalState()); @@ -154,8 +157,7 @@ public class ExecutionWithContainer { try { logger.info("Execution object: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject)); } catch (JsonProcessingException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + throw new RuntimeException(e); } } @@ -170,8 +172,7 @@ public class ExecutionWithContainer { try { logger.info("Execution object in CREATED State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject)); } catch (JsonProcessingException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + throw new RuntimeException(e); } for (Step s : execObject.getContext().getSteps().values()) { if (s.getInputs().stream().anyMatch(i -> i.getDescriptor().getName().equals("name") && !i.isRegistered())) { @@ -184,8 +185,7 @@ public class ExecutionWithContainer { try { logger.info("Execution object in READY State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject)); } catch (JsonProcessingException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + throw new RuntimeException(e); } @@ -196,7 +196,8 @@ public class ExecutionWithContainer { logger.debug("Execution status: " + execObject.getContext().getStatus()); execObject = executionsService.getExecution(execObject.getId()); } catch (InterruptedException e) { - e.printStackTrace(); + Thread.currentThread().interrupt(); + throw new RuntimeException(e); } } assertEquals(ExecutionStatus.WAITING, execObject.getContext().getStatus()); @@ -204,8 +205,7 @@ public class ExecutionWithContainer { try { logger.info("Execution object in WAITING State: {}", ObjectMapperHolder.mapper.writerWithDefaultPrettyPrinter().writeValueAsString(execObject)); } catch (JsonProcessingException e) { - // TODO Auto-generated catch block - e.printStackTrace(); + throw new RuntimeException(e); } }