diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockDraftNormalizer.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockDraftNormalizer.java new file mode 100644 index 0000000..42482e0 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/BlockDraftNormalizer.java @@ -0,0 +1,418 @@ +package it.cnr.isti.workflow.manager.assistant; + +import java.util.List; +import java.util.Locale; +import java.util.Objects; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.http.HttpStatus; +import org.springframework.web.server.ResponseStatusException; + +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.node.ArrayNode; +import tools.jackson.databind.node.ObjectNode; + +import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantBlockPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantConfiguredBlockDraft; +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.factories.BlockFactory; +import it.cnr.isti.workflow.manager.mcp.MCPServersProvider; + +final class BlockDraftNormalizer { + + private static final Logger log = LoggerFactory.getLogger(BlockDraftNormalizer.class); + + private BlockDraftNormalizer() { + } + + static AssistantConfiguredBlockDraft normalizeBlockDraft(AssistantBlockPlan blockPlan, + AssistantConfiguredBlockDraft draft) { + if (draft == null) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an empty block configuration payload"); + } + if (draft.config() == null || draft.config().isMissingNode() || draft.config().isNull()) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a block configuration without config"); + } + if (draft.blockId() == null || draft.blockId().isBlank() || !Objects.equals(blockPlan.blockId(), draft.blockId())) { + return new AssistantConfiguredBlockDraft(blockPlan.blockId(), draft.name(), draft.config()); + } + return draft; + } + + static Block buildBlock(BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan, + AssistantConfiguredBlockDraft draft, String provider, String model, boolean requireSharedMemorySemantics, + int blockIndex, int blockCount, MCPServersProvider mcpServersProvider, List> blockFactories) { + if (!(draft.config() instanceof ObjectNode configNode)) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a non-object config for block " + blockPlan.blockId()); + } + + ObjectNode normalizedConfig = configNode.deepCopy(); + normalizedConfig.put("type", descriptor.configurationType()); + normalizedConfig.put("name", AssistantTextSupport.defaultIfBlank(draft.name(), AssistantTextSupport.defaultIfBlank(blockPlan.purpose(), blockPlan.blockType()))); + injectSystemManagedFields(normalizedConfig, descriptor, provider, model); + ensureRequiredTextDefaults(normalizedConfig, descriptor, blockPlan); + normalizeHumanDecisionOptions(normalizedConfig, descriptor); + normalizeHttpServerCallAuthorization(normalizedConfig, blockPlan); + normalizeMcpAgentServers(normalizedConfig, blockPlan, mcpServersProvider); + normalizeMcpAgentSharedMemory(normalizedConfig, blockPlan, model, requireSharedMemorySemantics); + ensureMcpAgentModelConfigured(normalizedConfig, blockPlan, model); + ensureSequentialInputPlaceholder(normalizedConfig, descriptor, blockPlan, blockIndex, blockCount); + + try { + BlockConfiguration configuration = ObjectMapperHolder.mapper.treeToValue(normalizedConfig, + BlockConfiguration.class); + return createBlock(configuration, blockFactories); + } catch (Exception e) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an invalid block configuration for " + blockPlan.blockType() + ": " + + e.getMessage()); + } + } + + private static void normalizeHttpServerCallAuthorization(ObjectNode config, AssistantBlockPlan blockPlan) { + if (!"HTTPServerCall".equals(blockPlan.blockType())) { + return; + } + + JsonNode requiresAuthorizationNode = config.get("requiresAuthorization"); + boolean requiresAuthorization = requiresAuthorizationNode != null + && !requiresAuthorizationNode.isNull() + && requiresAuthorizationNode.asBoolean(false); + + JsonNode authorizationTypeNode = config.get("authorizationType"); + String authorizationType = authorizationTypeNode == null || authorizationTypeNode.isNull() + ? null + : authorizationTypeNode.stringValueOpt().orElse(null); + + if (authorizationType == null || authorizationType.isBlank()) { + if (!requiresAuthorization) { + config.put("requiresAuthorization", false); + config.put("authorizationType", "API_KEY"); + } + return; + } + + if ("NONE".equalsIgnoreCase(authorizationType)) { + config.put("requiresAuthorization", false); + config.put("authorizationType", "API_KEY"); + } + } + + private static void ensureSequentialInputPlaceholder(ObjectNode config, + BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan, int blockIndex, + int blockCount) { + if (blockIndex <= 0 || blockCount < 2 || !isPromptDrivenConfiguration(descriptor.configurationType())) { + return; + } + + String prompt = AssistantTextSupport.textOrEmpty(config.path("prompt")); + if (prompt.contains("${{")) { + return; + } + + String dependencyHint = "Use upstream workflow context from ${{input}}."; + String normalizedPrompt = prompt.isBlank() + ? dependencyHint + : prompt.stripTrailing() + "\n\n" + dependencyHint; + config.put("prompt", normalizedPrompt); + log.debug("Added default upstream input placeholder to assistant-configured block {} ({})", + blockPlan.blockId(), + blockPlan.blockType()); + } + + private static boolean isPromptDrivenConfiguration(String configurationType) { + return isConfigurationType(configurationType, "LLMBlockConfiguration") + || isConfigurationType(configurationType, "MCPAgentBlockConfiguration"); + } + + /** + * Drops CATALOG-sourced MCP server bindings whose serverName is not present in the declared catalog, + * so a model that hallucinated a server id degrades to a clean validation/fallback path instead of a + * runtime "Unknown MCP server" failure. CUSTOM bindings (self-described servers) are left untouched. + */ + private static void normalizeMcpAgentServers(ObjectNode config, AssistantBlockPlan blockPlan, + MCPServersProvider mcpServersProvider) { + String blockType = blockPlan.blockType(); + if (!"MCPAgent".equals(blockType) && !"MCPAgentChat".equals(blockType)) { + return; + } + JsonNode serversNode = config.get("mcpServers"); + if (serversNode == null || !serversNode.isArray()) { + return; + } + ArrayNode kept = ObjectMapperHolder.mapper.createArrayNode(); + for (JsonNode server : serversNode) { + String sourceType = AssistantTextSupport.textOrEmpty(server.path("sourceType")); + boolean catalogSourced = sourceType.isBlank() || "CATALOG".equalsIgnoreCase(sourceType); + if (!catalogSourced) { + kept.add(server); + continue; + } + String serverName = AssistantTextSupport.textOrEmpty(server.path("serverName")); + if (isKnownMcpServer(serverName, mcpServersProvider)) { + kept.add(server); + } else { + log.warn("Dropping unknown MCP server '{}' chosen for block {} ({})", + serverName, blockPlan.blockId(), blockType); + } + } + if (kept.isEmpty()) { + config.remove("mcpServers"); + } else { + config.set("mcpServers", kept); + } + } + + private static void normalizeMcpAgentSharedMemory(ObjectNode config, AssistantBlockPlan blockPlan, String model, + boolean requireSharedMemorySemantics) { + if (!requireSharedMemorySemantics || !"MCPAgent".equals(blockPlan.blockType())) { + return; + } + + boolean producerPurpose = SharedMemoryIntentClassifier.isSharedStateProducerPurpose(blockPlan.purpose()); + boolean consumerPurpose = SharedMemoryIntentClassifier.isSharedStateConsumerPurpose(blockPlan.purpose()); + + if (producerPurpose && (!consumerPurpose || SharedMemoryIntentClassifier.isSharedStateProducerDominantPurpose(blockPlan.purpose()))) { + configureMcpSharedMemoryProducer(config, model); + return; + } + + if (consumerPurpose) { + configureMcpSharedMemoryConsumer(config); + return; + } + + if (producerPurpose) { + configureMcpSharedMemoryProducer(config, model); + } + } + + private static void configureMcpSharedMemoryProducer(ObjectNode config, String model) { + config.put("shareSession", true); + config.put("useSharedSession", false); + if (!AssistantTextSupport.hasTextValue(config.get("sharedSessionName"))) { + config.put("sharedSessionName", FlowAssistantService.SHARED_MEMORY_SESSION_NAME); + } + if (!AssistantTextSupport.hasTextValue(config.get("model"))) { + config.put("model", model); + } + ensureDefaultRagServer(config); + } + + private static void configureMcpSharedMemoryConsumer(ObjectNode config) { + config.put("shareSession", false); + config.put("useSharedSession", true); + if (!AssistantTextSupport.hasTextValue(config.get("sharedSessionRef"))) { + config.put("sharedSessionRef", FlowAssistantService.SHARED_MEMORY_SESSION_NAME); + } + } + + /** + * MCPAgent/MCPAgentChat expose {@code model} as a {@code @ConfigurableAsInput} field: the block + * factory turns it into a real block INPUT whenever the config leaves it blank. An assistant + * that omits the model therefore produces a stray "model" input that leaks into a container's + * exposed interface (and, for a LoopContainer, counts toward the feedbackInput-ambiguity check, + * making the flow non-executable). model is not flow data - it is the agent's LLM - so fill it + * with the workflow model when the assistant left it blank, eliminating the phantom input. + */ + private static void ensureMcpAgentModelConfigured(ObjectNode config, AssistantBlockPlan blockPlan, String model) { + String blockType = blockPlan.blockType(); + if (!"MCPAgent".equals(blockType) && !"MCPAgentChat".equals(blockType)) { + return; + } + if (!AssistantTextSupport.hasTextValue(config.get("model"))) { + config.put("model", model); + } + } + + private static void ensureDefaultRagServer(ObjectNode config) { + if (hasMcpServer(config, "rag")) { + return; + } + JsonNode existingServers = config.get("mcpServers"); + if (existingServers != null && existingServers.isArray() && !existingServers.isEmpty()) { + return; + } + ArrayNode servers = config.putArray("mcpServers"); + ObjectNode ragServer = ObjectMapperHolder.mapper.createObjectNode(); + ragServer.put("sourceType", "CATALOG"); + ragServer.put("serverName", "rag"); + ragServer.set("configuration", ObjectMapperHolder.mapper.createObjectNode()); + servers.add(ragServer); + } + + static List mcpServerCatalogEntries(MCPServersProvider mcpServersProvider) { + return mcpServersProvider.getServers().stream() + .map(server -> new FlowAssistantPromptService.McpServerCatalogEntry( + server.id(), server.name(), server.description())) + .toList(); + } + + private static boolean isKnownMcpServer(String serverName, MCPServersProvider mcpServersProvider) { + if (serverName == null || serverName.isBlank()) { + return false; + } + return mcpServersProvider.getServers().stream() + .anyMatch(server -> serverName.equalsIgnoreCase(server.id()) + || serverName.equalsIgnoreCase(server.name())); + } + + private static boolean hasMcpServer(ObjectNode config, String serverName) { + JsonNode servers = config.get("mcpServers"); + if (servers == null || !servers.isArray()) { + return false; + } + for (JsonNode server : servers) { + if (serverName.equalsIgnoreCase(AssistantTextSupport.textOrEmpty(server.path("serverName")))) { + return true; + } + } + return false; + } + + private static void injectSystemManagedFields(ObjectNode config, BlockCatalogService.AssistantPromptBlockDescriptor descriptor, + String provider, String model) { + String configurationType = descriptor.configurationType(); + removeSystemManagedFields(config, descriptor); + if (isConfigurationType(configurationType, "LLMBlockConfiguration") + || isConfigurationType(configurationType, "ChatInteractionBlockConfiguration")) { + config.set("llmDescriptor", llmDescriptorNode(provider, model)); + return; + } + if (isConfigurationType(configurationType, "ConditionalBlockConfiguration") + || isConfigurationType(configurationType, "SwitchBlockConfiguration")) { + boolean useLlm = inferConditionalUseLlm(config); + config.put("useLlm", useLlm); + if (useLlm) { + config.set("llmDescriptor", llmDescriptorNode(provider, model)); + } + } + } + + /** + * Fills any required free-text configuration field the model omitted with a sensible default + * derived from the block's purpose, so a missing required string (e.g. MCPAgentChat's + * goalDescription) doesn't fail deserialization with a hard 502. Only plain required string + * fields are defaulted - enum fields (with allowed values) and structural fields are left + * untouched, as is any field the model already set. + */ + private static void ensureRequiredTextDefaults(ObjectNode config, + BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan) { + if (descriptor.configurationFields() == null) { + return; + } + for (BlockCatalogService.AssistantPromptFieldDescriptor field : descriptor.configurationFields()) { + if (!field.required() || field.structural() || !"string".equalsIgnoreCase(field.type())) { + continue; + } + if (field.allowedValues() != null && !field.allowedValues().isEmpty()) { + continue; + } + if (!AssistantTextSupport.hasTextValue(config.get(field.name()))) { + config.put(field.name(), AssistantTextSupport.defaultIfBlank(blockPlan.purpose(), "Assist the user with this step.")); + } + } + } + + /** + * HumanDecisionOption is (name, label): name is the routing key / branch output, label the + * display text. Models routinely emit "value" (or only "label") instead of "name", leaving name + * null - which produces null-named branch outputs. Fill a missing option name from its "value" + * (the intended routing key) or from a slug of its "label", so the flow is valid and each branch + * has a real output name. + */ + private static void normalizeHumanDecisionOptions(ObjectNode config, + BlockCatalogService.AssistantPromptBlockDescriptor descriptor) { + if (!isConfigurationType(descriptor.configurationType(), "HumanDecisionBlockConfiguration")) { + return; + } + if (!(config.get("options") instanceof ArrayNode options)) { + return; + } + for (JsonNode option : options) { + if (!(option instanceof ObjectNode optionObject) || AssistantTextSupport.hasTextValue(optionObject.get("name"))) { + continue; + } + String derived = AssistantTextSupport.hasTextValue(optionObject.get("value")) + ? AssistantTextSupport.textOrNull(optionObject.get("value")).trim() + : AssistantTextSupport.hasTextValue(optionObject.get("label")) + ? slugifyOptionName(AssistantTextSupport.textOrNull(optionObject.get("label"))) + : null; + if (derived != null && !derived.isBlank()) { + optionObject.put("name", derived); + } + } + } + + private static String slugifyOptionName(String label) { + String slug = label.trim().toLowerCase(Locale.ROOT).replaceAll("[^a-z0-9]+", "-").replaceAll("(^-+|-+$)", ""); + if (slug.isBlank()) { + return "option"; + } + return Character.isLetter(slug.charAt(0)) ? slug : "opt-" + slug; + } + + /** + * Drops the fields the backend owns so a model echoing them cannot overwrite system state - + * except a field the prompt catalog exposed to the assistant as both structural and required. + * That is exactly the IO the assistant was asked to declare and without which the configuration + * cannot be built at all: BranchRejoinBlock's branch "inputs" and DelimitedParserBlock's + * "outputs". Stripping those deleted a correct model answer right before deserialization and + * surfaced as a missing required-property failure blamed on the assistant - one the repair loop + * could never fix, since every repaired response was stripped again. Structural-but-optional IO + * (e.g. ChatInteraction's inputs, derived from the prompt placeholders) stays system-managed. + */ + private static void removeSystemManagedFields(ObjectNode config, + BlockCatalogService.AssistantPromptBlockDescriptor descriptor) { + config.remove(FlowAssistantService.SYSTEM_MANAGED_FIELDS.stream() + .filter(field -> !isAssistantDeclaredField(descriptor, field)) + .toList()); + } + + private static boolean isAssistantDeclaredField(BlockCatalogService.AssistantPromptBlockDescriptor descriptor, + String fieldName) { + if (descriptor == null || descriptor.configurationFields() == null) { + return false; + } + return descriptor.configurationFields().stream() + .anyMatch(field -> field.structural() && field.required() && fieldName.equals(field.name())); + } + + private static boolean isConfigurationType(String actualType, String expectedSimpleName) { + return actualType != null + && (actualType.equals(expectedSimpleName) || actualType.endsWith("." + expectedSimpleName)); + } + + private static boolean inferConditionalUseLlm(ObjectNode config) { + if (config.has("useLlm")) { + return config.get("useLlm").asBoolean(false); + } + if (config.hasNonNull("prompt")) { + return true; + } + return false; + } + + @SuppressWarnings({ "rawtypes", "unchecked" }) + private static Block createBlock(BlockConfiguration configuration, List> blockFactories) { + BlockFactory factory = blockFactories.stream() + .filter(candidate -> candidate.getBlockType().equals(configuration.getBlockType())) + .findFirst() + .orElseThrow(() -> new IllegalArgumentException( + "Block factory not found for type: " + configuration.getBlockType().getSimpleName())); + return (Block) factory.create(configuration); + } + + private static ObjectNode llmDescriptorNode(String provider, String model) { + ObjectNode llmDescriptor = ObjectMapperHolder.mapper.createObjectNode(); + llmDescriptor.put("provider", provider); + llmDescriptor.put("model", model); + return llmDescriptor; + } +} 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 9ceaf0c..0c87dc9 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 @@ -23,8 +23,6 @@ import org.springframework.web.server.ResponseStatusException; import tools.jackson.databind.JsonNode; import tools.jackson.databind.ObjectMapper; -import tools.jackson.databind.node.ArrayNode; -import tools.jackson.databind.node.ObjectNode; import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; import it.cnr.isti.workflow.manager.assistant.FlowAssistantPromptService.OperationMode; @@ -37,7 +35,6 @@ import it.cnr.isti.workflow.manager.assistant.model.AssistantLlmSelection; import it.cnr.isti.workflow.manager.assistant.model.AssistantModelSelection; import it.cnr.isti.workflow.manager.assistant.model.AssistantRefineRequest; import it.cnr.isti.workflow.manager.blocks.Block; -import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfiguration; @@ -80,13 +77,13 @@ public class FlowAssistantService { private static final Logger assistantResponseLog = LoggerFactory.getLogger("assistant.responses"); private static final String INTERNAL_PROVIDER_NAME = "InternalOllama"; - private static final String SHARED_MEMORY_SESSION_NAME = "sharedMemorySession"; + static final String SHARED_MEMORY_SESSION_NAME = "sharedMemorySession"; private static final int DEFAULT_PROVIDER_RETRY_ATTEMPTS = 3; private static final int DEFAULT_MAX_REPAIR_ATTEMPTS = 2; // Fields the backend owns: the assistant never gets to set them, and anything it echoes back // is dropped before deserialization. "inputs"/"outputs" are in here because they are usually // runtime-derived IO lists - see removeSystemManagedFields for the structural exception. - private static final List SYSTEM_MANAGED_FIELDS = List.of("provider", "model", "llmDescriptor", "ids", + static final List SYSTEM_MANAGED_FIELDS = List.of("provider", "model", "llmDescriptor", "ids", "inputs", "outputs", "skills"); // Error codes that imply the flow's shape itself (block set, block I/O, connections, // dependencies, lanes, shared-session wiring) may need to change. Any of these - or any @@ -470,7 +467,7 @@ public class FlowAssistantService { parsedPlan = PlanValidationSupport.buildReusedPlanForTargetedRepair(currentFlow, errors); } else { String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog, - mcpServerCatalogEntries()); + BlockDraftNormalizer.mcpServerCatalogEntries(mcpServersProvider)); parsedPlan = invokeStructuredAndValidate(provider, authorization, planningModelFor(mode, phaseModels), phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { ParsedPlan plan = AssistantResponseParser.parsePlan(rawResponse); @@ -522,17 +519,17 @@ public class FlowAssistantService { } else { progressListener.onProgress("configuring_blocks", "Configuring block " + blockPlan.blockId()); String blockPrompt = promptService.buildBlockConfigurationPrompt(mode, userPrompt, descriptor, - parsedPlan.plan(), blockPlan, currentFlow, errors, assistantModel, mcpServerCatalogEntries()); + parsedPlan.plan(), blockPlan, currentFlow, errors, assistantModel, BlockDraftNormalizer.mcpServerCatalogEntries(mcpServersProvider)); int currentBlockIndex = blockIndex; int blockPlanCount = blockPlans.size(); ConfiguredBlockResult configuredBlock = invokeStructuredAndValidate(provider, authorization, jsonModelFor(mode, phaseModels), phaseModels.repairModel(), blockPrompt, "block configuration for " + blockPlan.blockId(), rawResponse -> { ParsedBlockDraft parsedBlock = AssistantResponseParser.parseBlockDraft(rawResponse); - AssistantConfiguredBlockDraft normalizedDraft = normalizeBlockDraft(blockPlan, + AssistantConfiguredBlockDraft normalizedDraft = BlockDraftNormalizer.normalizeBlockDraft(blockPlan, parsedBlock.block()); - Block newBlock = buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, - requireSharedMemorySemantics, currentBlockIndex, blockPlanCount); + Block newBlock = BlockDraftNormalizer.buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, + requireSharedMemorySemantics, currentBlockIndex, blockPlanCount, mcpServersProvider, blockFactories); return new ConfiguredBlockResult(parsedBlock, newBlock); }); appendRationale(rationaleParts, configuredBlock.parsedBlock().rationale()); @@ -723,7 +720,7 @@ public class FlowAssistantService { "Configuring block " + blockPlan.blockId() + " in container " + containerPlan.containerId()); String blockPrompt = promptService.buildBlockConfigurationPrompt(mode, userPrompt, descriptor, containerInnerPlan(containerPlan), blockPlan, null, List.of(), assistantModel, - mcpServerCatalogEntries()); + BlockDraftNormalizer.mcpServerCatalogEntries(mcpServersProvider)); int currentBlockIndex = blockIndex; int blockPlanCount = innerBlockPlans.size(); ConfiguredBlockResult configuredBlock = invokeStructuredAndValidate(provider, authorization, @@ -731,9 +728,9 @@ public class FlowAssistantService { "block configuration for " + blockPlan.blockId() + " in container " + containerPlan.containerId(), rawResponse -> { ParsedBlockDraft parsedBlock = AssistantResponseParser.parseBlockDraft(rawResponse); - AssistantConfiguredBlockDraft normalizedDraft = normalizeBlockDraft(blockPlan, parsedBlock.block()); - Block newBlock = buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, - innerRequiresSharedMemory, currentBlockIndex, blockPlanCount); + AssistantConfiguredBlockDraft normalizedDraft = BlockDraftNormalizer.normalizeBlockDraft(blockPlan, parsedBlock.block()); + Block newBlock = BlockDraftNormalizer.buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, + innerRequiresSharedMemory, currentBlockIndex, blockPlanCount, mcpServersProvider, blockFactories); return new ConfiguredBlockResult(parsedBlock, newBlock); }); appendRationale(rationaleParts, configuredBlock.parsedBlock().rationale()); @@ -1119,386 +1116,6 @@ public class FlowAssistantService { } } - private AssistantConfiguredBlockDraft normalizeBlockDraft(AssistantBlockPlan blockPlan, - AssistantConfiguredBlockDraft draft) { - if (draft == null) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an empty block configuration payload"); - } - if (draft.config() == null || draft.config().isMissingNode() || draft.config().isNull()) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a block configuration without config"); - } - if (draft.blockId() == null || draft.blockId().isBlank() || !Objects.equals(blockPlan.blockId(), draft.blockId())) { - return new AssistantConfiguredBlockDraft(blockPlan.blockId(), draft.name(), draft.config()); - } - return draft; - } - - private Block buildBlock(BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan, - AssistantConfiguredBlockDraft draft, String provider, String model, boolean requireSharedMemorySemantics, int blockIndex, - int blockCount) { - if (!(draft.config() instanceof ObjectNode configNode)) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a non-object config for block " + blockPlan.blockId()); - } - - ObjectNode normalizedConfig = configNode.deepCopy(); - normalizedConfig.put("type", descriptor.configurationType()); - normalizedConfig.put("name", AssistantTextSupport.defaultIfBlank(draft.name(), AssistantTextSupport.defaultIfBlank(blockPlan.purpose(), blockPlan.blockType()))); - injectSystemManagedFields(normalizedConfig, descriptor, provider, model); - ensureRequiredTextDefaults(normalizedConfig, descriptor, blockPlan); - normalizeHumanDecisionOptions(normalizedConfig, descriptor); - normalizeHttpServerCallAuthorization(normalizedConfig, blockPlan); - normalizeMcpAgentServers(normalizedConfig, blockPlan); - normalizeMcpAgentSharedMemory(normalizedConfig, blockPlan, model, requireSharedMemorySemantics); - ensureMcpAgentModelConfigured(normalizedConfig, blockPlan, model); - ensureSequentialInputPlaceholder(normalizedConfig, descriptor, blockPlan, blockIndex, blockCount); - - try { - BlockConfiguration configuration = ObjectMapperHolder.mapper.treeToValue(normalizedConfig, - BlockConfiguration.class); - return createBlock(configuration); - } catch (Exception e) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an invalid block configuration for " + blockPlan.blockType() + ": " - + e.getMessage()); - } - } - - private void normalizeHttpServerCallAuthorization(ObjectNode config, AssistantBlockPlan blockPlan) { - if (!"HTTPServerCall".equals(blockPlan.blockType())) { - return; - } - - JsonNode requiresAuthorizationNode = config.get("requiresAuthorization"); - boolean requiresAuthorization = requiresAuthorizationNode != null - && !requiresAuthorizationNode.isNull() - && requiresAuthorizationNode.asBoolean(false); - - JsonNode authorizationTypeNode = config.get("authorizationType"); - String authorizationType = authorizationTypeNode == null || authorizationTypeNode.isNull() - ? null - : authorizationTypeNode.stringValueOpt().orElse(null); - - if (authorizationType == null || authorizationType.isBlank()) { - if (!requiresAuthorization) { - config.put("requiresAuthorization", false); - config.put("authorizationType", "API_KEY"); - } - return; - } - - if ("NONE".equalsIgnoreCase(authorizationType)) { - config.put("requiresAuthorization", false); - config.put("authorizationType", "API_KEY"); - } - } - - private void ensureSequentialInputPlaceholder(ObjectNode config, - BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan, int blockIndex, - int blockCount) { - if (blockIndex <= 0 || blockCount < 2 || !isPromptDrivenConfiguration(descriptor.configurationType())) { - return; - } - - String prompt = AssistantTextSupport.textOrEmpty(config.path("prompt")); - if (prompt.contains("${{")) { - return; - } - - String dependencyHint = "Use upstream workflow context from ${{input}}."; - String normalizedPrompt = prompt.isBlank() - ? dependencyHint - : prompt.stripTrailing() + "\n\n" + dependencyHint; - config.put("prompt", normalizedPrompt); - log.debug("Added default upstream input placeholder to assistant-configured block {} ({})", - blockPlan.blockId(), - blockPlan.blockType()); - } - - private boolean isPromptDrivenConfiguration(String configurationType) { - return isConfigurationType(configurationType, "LLMBlockConfiguration") - || isConfigurationType(configurationType, "MCPAgentBlockConfiguration"); - } - - /** - * Drops CATALOG-sourced MCP server bindings whose serverName is not present in the declared catalog, - * so a model that hallucinated a server id degrades to a clean validation/fallback path instead of a - * runtime "Unknown MCP server" failure. CUSTOM bindings (self-described servers) are left untouched. - */ - private void normalizeMcpAgentServers(ObjectNode config, AssistantBlockPlan blockPlan) { - String blockType = blockPlan.blockType(); - if (!"MCPAgent".equals(blockType) && !"MCPAgentChat".equals(blockType)) { - return; - } - JsonNode serversNode = config.get("mcpServers"); - if (serversNode == null || !serversNode.isArray()) { - return; - } - ArrayNode kept = ObjectMapperHolder.mapper.createArrayNode(); - for (JsonNode server : serversNode) { - String sourceType = AssistantTextSupport.textOrEmpty(server.path("sourceType")); - boolean catalogSourced = sourceType.isBlank() || "CATALOG".equalsIgnoreCase(sourceType); - if (!catalogSourced) { - kept.add(server); - continue; - } - String serverName = AssistantTextSupport.textOrEmpty(server.path("serverName")); - if (isKnownMcpServer(serverName)) { - kept.add(server); - } else { - log.warn("Dropping unknown MCP server '{}' chosen for block {} ({})", - serverName, blockPlan.blockId(), blockType); - } - } - if (kept.isEmpty()) { - config.remove("mcpServers"); - } else { - config.set("mcpServers", kept); - } - } - - private void normalizeMcpAgentSharedMemory(ObjectNode config, AssistantBlockPlan blockPlan, String model, - boolean requireSharedMemorySemantics) { - if (!requireSharedMemorySemantics || !"MCPAgent".equals(blockPlan.blockType())) { - return; - } - - boolean producerPurpose = SharedMemoryIntentClassifier.isSharedStateProducerPurpose(blockPlan.purpose()); - boolean consumerPurpose = SharedMemoryIntentClassifier.isSharedStateConsumerPurpose(blockPlan.purpose()); - - if (producerPurpose && (!consumerPurpose || SharedMemoryIntentClassifier.isSharedStateProducerDominantPurpose(blockPlan.purpose()))) { - configureMcpSharedMemoryProducer(config, model); - return; - } - - if (consumerPurpose) { - configureMcpSharedMemoryConsumer(config); - return; - } - - if (producerPurpose) { - configureMcpSharedMemoryProducer(config, model); - } - } - - private void configureMcpSharedMemoryProducer(ObjectNode config, String model) { - config.put("shareSession", true); - config.put("useSharedSession", false); - if (!AssistantTextSupport.hasTextValue(config.get("sharedSessionName"))) { - config.put("sharedSessionName", SHARED_MEMORY_SESSION_NAME); - } - if (!AssistantTextSupport.hasTextValue(config.get("model"))) { - config.put("model", model); - } - ensureDefaultRagServer(config); - } - - private void configureMcpSharedMemoryConsumer(ObjectNode config) { - config.put("shareSession", false); - config.put("useSharedSession", true); - if (!AssistantTextSupport.hasTextValue(config.get("sharedSessionRef"))) { - config.put("sharedSessionRef", SHARED_MEMORY_SESSION_NAME); - } - } - - /** - * MCPAgent/MCPAgentChat expose {@code model} as a {@code @ConfigurableAsInput} field: the block - * factory turns it into a real block INPUT whenever the config leaves it blank. An assistant - * that omits the model therefore produces a stray "model" input that leaks into a container's - * exposed interface (and, for a LoopContainer, counts toward the feedbackInput-ambiguity check, - * making the flow non-executable). model is not flow data - it is the agent's LLM - so fill it - * with the workflow model when the assistant left it blank, eliminating the phantom input. - */ - private void ensureMcpAgentModelConfigured(ObjectNode config, AssistantBlockPlan blockPlan, String model) { - String blockType = blockPlan.blockType(); - if (!"MCPAgent".equals(blockType) && !"MCPAgentChat".equals(blockType)) { - return; - } - if (!AssistantTextSupport.hasTextValue(config.get("model"))) { - config.put("model", model); - } - } - - private void ensureDefaultRagServer(ObjectNode config) { - if (hasMcpServer(config, "rag")) { - return; - } - JsonNode existingServers = config.get("mcpServers"); - if (existingServers != null && existingServers.isArray() && !existingServers.isEmpty()) { - return; - } - ArrayNode servers = config.putArray("mcpServers"); - ObjectNode ragServer = ObjectMapperHolder.mapper.createObjectNode(); - ragServer.put("sourceType", "CATALOG"); - ragServer.put("serverName", "rag"); - ragServer.set("configuration", ObjectMapperHolder.mapper.createObjectNode()); - servers.add(ragServer); - } - - private List mcpServerCatalogEntries() { - return mcpServersProvider.getServers().stream() - .map(server -> new FlowAssistantPromptService.McpServerCatalogEntry( - server.id(), server.name(), server.description())) - .toList(); - } - - private boolean isKnownMcpServer(String serverName) { - if (serverName == null || serverName.isBlank()) { - return false; - } - return mcpServersProvider.getServers().stream() - .anyMatch(server -> serverName.equalsIgnoreCase(server.id()) - || serverName.equalsIgnoreCase(server.name())); - } - - private boolean hasMcpServer(ObjectNode config, String serverName) { - JsonNode servers = config.get("mcpServers"); - if (servers == null || !servers.isArray()) { - return false; - } - for (JsonNode server : servers) { - if (serverName.equalsIgnoreCase(AssistantTextSupport.textOrEmpty(server.path("serverName")))) { - return true; - } - } - return false; - } - - private void injectSystemManagedFields(ObjectNode config, BlockCatalogService.AssistantPromptBlockDescriptor descriptor, - String provider, String model) { - String configurationType = descriptor.configurationType(); - removeSystemManagedFields(config, descriptor); - if (isConfigurationType(configurationType, "LLMBlockConfiguration") - || isConfigurationType(configurationType, "ChatInteractionBlockConfiguration")) { - config.set("llmDescriptor", llmDescriptorNode(provider, model)); - return; - } - if (isConfigurationType(configurationType, "ConditionalBlockConfiguration") - || isConfigurationType(configurationType, "SwitchBlockConfiguration")) { - boolean useLlm = inferConditionalUseLlm(config); - config.put("useLlm", useLlm); - if (useLlm) { - config.set("llmDescriptor", llmDescriptorNode(provider, model)); - } - } - } - - /** - * Fills any required free-text configuration field the model omitted with a sensible default - * derived from the block's purpose, so a missing required string (e.g. MCPAgentChat's - * goalDescription) doesn't fail deserialization with a hard 502. Only plain required string - * fields are defaulted - enum fields (with allowed values) and structural fields are left - * untouched, as is any field the model already set. - */ - private void ensureRequiredTextDefaults(ObjectNode config, - BlockCatalogService.AssistantPromptBlockDescriptor descriptor, AssistantBlockPlan blockPlan) { - if (descriptor.configurationFields() == null) { - return; - } - for (BlockCatalogService.AssistantPromptFieldDescriptor field : descriptor.configurationFields()) { - if (!field.required() || field.structural() || !"string".equalsIgnoreCase(field.type())) { - continue; - } - if (field.allowedValues() != null && !field.allowedValues().isEmpty()) { - continue; - } - if (!AssistantTextSupport.hasTextValue(config.get(field.name()))) { - config.put(field.name(), AssistantTextSupport.defaultIfBlank(blockPlan.purpose(), "Assist the user with this step.")); - } - } - } - - /** - * HumanDecisionOption is (name, label): name is the routing key / branch output, label the - * display text. Models routinely emit "value" (or only "label") instead of "name", leaving name - * null - which produces null-named branch outputs. Fill a missing option name from its "value" - * (the intended routing key) or from a slug of its "label", so the flow is valid and each branch - * has a real output name. - */ - private void normalizeHumanDecisionOptions(ObjectNode config, - BlockCatalogService.AssistantPromptBlockDescriptor descriptor) { - if (!isConfigurationType(descriptor.configurationType(), "HumanDecisionBlockConfiguration")) { - return; - } - if (!(config.get("options") instanceof ArrayNode options)) { - return; - } - for (JsonNode option : options) { - if (!(option instanceof ObjectNode optionObject) || AssistantTextSupport.hasTextValue(optionObject.get("name"))) { - continue; - } - String derived = AssistantTextSupport.hasTextValue(optionObject.get("value")) - ? AssistantTextSupport.textOrNull(optionObject.get("value")).trim() - : AssistantTextSupport.hasTextValue(optionObject.get("label")) - ? slugifyOptionName(AssistantTextSupport.textOrNull(optionObject.get("label"))) - : null; - if (derived != null && !derived.isBlank()) { - optionObject.put("name", derived); - } - } - } - - private String slugifyOptionName(String label) { - String slug = label.trim().toLowerCase(Locale.ROOT).replaceAll("[^a-z0-9]+", "-").replaceAll("(^-+|-+$)", ""); - if (slug.isBlank()) { - return "option"; - } - return Character.isLetter(slug.charAt(0)) ? slug : "opt-" + slug; - } - - /** - * Drops the fields the backend owns so a model echoing them cannot overwrite system state - - * except a field the prompt catalog exposed to the assistant as both structural and required. - * That is exactly the IO the assistant was asked to declare and without which the configuration - * cannot be built at all: BranchRejoinBlock's branch "inputs" and DelimitedParserBlock's - * "outputs". Stripping those deleted a correct model answer right before deserialization and - * surfaced as a missing required-property failure blamed on the assistant - one the repair loop - * could never fix, since every repaired response was stripped again. Structural-but-optional IO - * (e.g. ChatInteraction's inputs, derived from the prompt placeholders) stays system-managed. - */ - private void removeSystemManagedFields(ObjectNode config, - BlockCatalogService.AssistantPromptBlockDescriptor descriptor) { - config.remove(SYSTEM_MANAGED_FIELDS.stream() - .filter(field -> !isAssistantDeclaredField(descriptor, field)) - .toList()); - } - - private boolean isAssistantDeclaredField(BlockCatalogService.AssistantPromptBlockDescriptor descriptor, - String fieldName) { - if (descriptor == null || descriptor.configurationFields() == null) { - return false; - } - return descriptor.configurationFields().stream() - .anyMatch(field -> field.structural() && field.required() && fieldName.equals(field.name())); - } - - private boolean isConfigurationType(String actualType, String expectedSimpleName) { - return actualType != null - && (actualType.equals(expectedSimpleName) || actualType.endsWith("." + expectedSimpleName)); - } - - private boolean inferConditionalUseLlm(ObjectNode config) { - if (config.has("useLlm")) { - return config.get("useLlm").asBoolean(false); - } - if (config.hasNonNull("prompt")) { - return true; - } - return false; - } - - @SuppressWarnings({ "rawtypes", "unchecked" }) - private Block createBlock(BlockConfiguration configuration) { - BlockFactory factory = blockFactories.stream() - .filter(candidate -> candidate.getBlockType().equals(configuration.getBlockType())) - .findFirst() - .orElseThrow(() -> new IllegalArgumentException( - "Block factory not found for type: " + configuration.getBlockType().getSimpleName())); - return (Block) factory.create(configuration); - } - private List preserveCurrentConnections(FlowCreateRequest currentFlow, Map oldNodeIdToAssembledNode, Set removedExistingNodeIds) { List sourceConnections = currentFlow == null || currentFlow.flow() == null @@ -1579,13 +1196,6 @@ public class FlowAssistantService { } - private ObjectNode llmDescriptorNode(String provider, String model) { - ObjectNode llmDescriptor = ObjectMapperHolder.mapper.createObjectNode(); - llmDescriptor.put("provider", provider); - llmDescriptor.put("model", model); - return llmDescriptor; - } - private void appendRationale(List target, String rationale) { if (rationale != null && !rationale.isBlank()) { target.add(rationale.trim());