refactor(assistant): extract block-config sanitization into BlockDraftNormalizer

Cluster G from the structural analysis: strips system-managed fields,
fills required defaults, injects llmDescriptor, normalizes MCP server
bindings and shared-memory wiring, HumanDecision option names, HTTP
authorization defaults - pure ObjectNode manipulation depending only on
mcpServersProvider and blockFactories (now passed as parameters instead
of instance fields).

- New BlockDraftNormalizer holds buildBlock/normalizeBlockDraft and the
  full sanitization tree: injectSystemManagedFields, ensureRequiredTextDefaults,
  normalizeHumanDecisionOptions, normalizeHttpServerCallAuthorization,
  normalizeMcpAgentServers, normalizeMcpAgentSharedMemory (+ producer/
  consumer configuration), ensureMcpAgentModelConfigured,
  ensureSequentialInputPlaceholder, createBlock, llmDescriptorNode, etc.
- SHARED_MEMORY_SESSION_NAME and SYSTEM_MANAGED_FIELDS promoted from
  private to package-private constants so the new class can reference
  them without duplication

2352 -> 1805 -> 1412 lines. 3350 -> 1412 total so far (-1938, ~58%).
Behavior-preserving: pure extraction plus explicit parameter-passing for
the two fields this cluster touches, no logic changes.

Verified with `mvn test`: 462 tests, 0 failures, 0 errors.

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-02 10:30:06 +02:00
parent 46a1fb2d3e
commit 55a13af364
2 changed files with 429 additions and 401 deletions

View File

@ -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<BlockFactory<?, ?>> 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<FlowAssistantPromptService.McpServerCatalogEntry> 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<BlockFactory<?, ?>> 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;
}
}

View File

@ -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<String> SYSTEM_MANAGED_FIELDS = List.of("provider", "model", "llmDescriptor", "ids",
static final List<String> 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<FlowAssistantPromptService.McpServerCatalogEntry> 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<Connection> preserveCurrentConnections(FlowCreateRequest currentFlow,
Map<String, FlowNode> oldNodeIdToAssembledNode, Set<String> removedExistingNodeIds) {
List<Connection> 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<String> target, String rationale) {
if (rationale != null && !rationale.isBlank()) {
target.add(rationale.trim());