From effb506572ec9f137c61a1c0c729cc4a9ae57b3b Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 3 Aug 2026 11:47:15 +0200 Subject: [PATCH] feat(assistant): make the MCP server catalog central to block selection The planning and block-config models never saw which MCP servers were declared, so they could not reach for tool-running steps (writing to a workspace, compiling, running shell) - the coding-agent-mcp server was effectively invisible and the assistant fell back to HTTPServerCall. - Inject the declared MCP catalog (id/name/description) into both the PLAN and BLOCK_CONFIG prompts, and steer MCPAgent/MCPAgentChat toward running declared-server tools instead of emulating them via HTTP. - Let the model bind a concrete catalog server via mcpServers; validate chosen serverName against the catalog and drop hallucinated ids so a bad choice degrades to a clean validation/fallback path instead of a runtime "Unknown MCP server" failure. The model's explicit choice is preserved and no longer overwritten by the rag default. Tests: catalog reaches both prompts; a chosen catalog server is kept; an unknown server id is dropped. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../assistant/FlowAssistantPromptService.java | 39 ++++++- .../assistant/FlowAssistantService.java | 65 ++++++++++- .../controllers/AssistantControllerTest.java | 103 ++++++++++++++++++ 3 files changed, 201 insertions(+), 6 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java index 686cc16..d5ba8a7 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java @@ -18,8 +18,16 @@ public class FlowAssistantPromptService { FIX } + /** + * Compact view of one declared MCP server, surfaced to the planning/config models so they can + * pick a concrete server (by id) for MCPAgent / MCPAgentChat steps instead of guessing. + */ + public record McpServerCatalogEntry(String id, String name, String description) { + } + public String buildPlanPrompt(OperationMode mode, String userPrompt, FlowCreateRequest currentFlow, - List errors, List catalog) { + List errors, List catalog, + List mcpServers) { String modeInstructions = switch (mode) { case REFINE -> """ MODE INSTRUCTIONS (REFINE): A flow already exists (see "Current flow" below). \ @@ -50,6 +58,7 @@ public class FlowAssistantPromptService { - Need user chat + no MCP memory/tools across nodes: prefer ChatInteraction. - Need user chat + MCP tools/shared session across nodes: prefer MCPAgentChat. - Need non-chat background processing with shared session/state: prefer MCPAgent. + - Need to run tools exposed by a declared MCP server - e.g. writing files to a workspace, building/compiling code, running shell commands, querying a database, downloading URLs, or indexing/RAG - prefer MCPAgent (no user in the loop) or MCPAgentChat (with a user in the loop), and later bind it to the concrete server from "Available MCP servers" whose description matches the task. Do NOT use HTTPServerCall to emulate a capability that a declared MCP server already provides (e.g. do not POST to an API to "save code" when an MCP server can write to the workspace directly). - Need non-chat stateless post-processing (including final answer validation): prefer LLMBlock. - If one step only checks or grades output from another step and does not need shared state, keep it as LLMBlock. - If a later step must reuse state produced earlier, use MCPAgent/MCPAgentChat for the stateful steps and connect producer to consumer explicitly. @@ -133,7 +142,7 @@ public class FlowAssistantPromptService { - Keep the plan minimal. - One block per logical action. - Prefer LLMBlock for stateless transformation, classification, summarization, extraction, validation, normalization, or synthesis steps. - - Prefer MCPAgent for indexing, retrieval, memory, persistent or shared workspace state, or any step that must preserve and reuse state produced by another block. + - Prefer MCPAgent when a step must run tools provided by a declared MCP server (see "Available MCP servers" below) - writing files to a workspace, building/compiling, running shell commands, querying a database, downloading URLs, indexing/retrieval/RAG - or must preserve and reuse state produced by another block. Only plan an MCPAgent step when the catalog actually contains a server whose tools fit the task. - Prefer HTTPServerCall for fetching files, APIs, or URLs before processing them. - If a task can be completed with a single stateless step, do not introduce MCPAgent. - If a later step must reuse data from an earlier step, make the earlier step produce it and the later step consume it explicitly. @@ -158,6 +167,9 @@ public class FlowAssistantPromptService { Available block catalog: %s + Available MCP servers (usable only by MCPAgent / MCPAgentChat blocks; pick one whose description matches a tool-running step): + %s + Current flow: %s @@ -171,6 +183,7 @@ public class FlowAssistantPromptService { modeInstructions, buildBlockChoiceGuide(catalog), toJson(catalog), + buildMcpServerCatalog(mcpServers), summarizeFlow(currentFlow), summarizeErrors(errors), userPrompt == null || userPrompt.isBlank() ? "(none)" : userPrompt); @@ -182,9 +195,24 @@ public class FlowAssistantPromptService { .collect(Collectors.joining("\n")); } + private String buildMcpServerCatalog(List mcpServers) { + if (mcpServers == null || mcpServers.isEmpty()) { + return "(no MCP servers are declared; do not plan MCPAgent/MCPAgentChat steps that need external tools)"; + } + return mcpServers.stream() + .map(server -> "- %s (%s): %s".formatted( + server.id(), + server.name() == null || server.name().isBlank() ? server.id() : server.name(), + server.description() == null || server.description().isBlank() + ? "(no description)" + : server.description())) + .collect(Collectors.joining("\n")); + } + public String buildBlockConfigurationPrompt(OperationMode mode, String userPrompt, BlockCatalogService.AssistantPromptBlockDescriptor descriptor, Object flowPlan, Object blockPlan, - FlowCreateRequest currentFlow, List errors, String selectedModel) { + FlowCreateRequest currentFlow, List errors, String selectedModel, + List mcpServers) { String modeInstructions = switch (mode) { case REFINE -> "MODE INSTRUCTIONS (REFINE): If this block already exists in the current flow and the user request does not ask to change it, preserve its existing configuration exactly. Only reconfigure it if the user request explicitly requires a change to this block."; case FIX -> "MODE INSTRUCTIONS (FIX): Reconfigure this block only to the extent needed to fix the reported validation errors. Keep everything else unchanged."; @@ -212,6 +240,7 @@ public class FlowAssistantPromptService { - For HTTPServerCall without authentication, set requiresAuthorization to false and omit authorization details. - MCPAgent shared session rule: when multiple blocks need to work with the SAME MCP tools AND share the same execution context or state, model them as a producer-consumer chain. The FIRST block that opens the tool session is the producer: set shareSession to true and sharedSessionName to "sharedMemorySession". Subsequent blocks that must reuse the same tools and the same live context (same session, same accumulated state) are consumers: set useSharedSession to true and sharedSessionRef to "sharedMemorySession". The backend orders the consumer after the producer automatically (via a dependency) - do NOT add any ordering placeholder or connection between them. A later block that does NOT need the same tools or the same context should NOT be made a consumer — use a separate independent MCPAgent or a different block type instead. - The producer is always the first block in the chain that initialises the shared tool session. Never mark it as a consumer. + - MCP server binding (MCPAgent / MCPAgentChat only): set "mcpServers" to the declared server(s) whose tools this step actually needs, each shaped {"sourceType": "CATALOG", "serverName": "", "configuration": {}} using an "id" listed in "Available MCP servers" below. Choose the server whose description matches what this step must do (e.g. a step that writes files or compiles/runs code -> the coding/workspace server; a step that runs SQL -> the database server; a retrieval/RAG step -> the RAG server). You may bind more than one server if the step needs tools from several. Do NOT invent server ids that are not in the list, and omit "mcpServers" only if the step genuinely needs no external tools. - A shared-session MCP producer/consumer chain may live all at the top level, all inside one container, or span the boundary from a TOP-LEVEL producer to a consumer inside a container - but in the last case the producer must be wired (with a connection) to run before that container, so its session is ready when the container runs. Do NOT put the producer inside a container and the consumer at the top level, and do NOT split a chain across two different containers. - Do not include system-managed fields like provider, model, llmDescriptor, ids, skills. The exception is BranchRejoinBlock's "inputs" field, which is a structural config array you must set (see below), not a derived IO list. - Use placeholders like ${{variable}} when needed. @@ -267,6 +296,9 @@ public class FlowAssistantPromptService { Block descriptor: %s + Available MCP servers (bind one or more via "mcpServers" only when configuring an MCPAgent / MCPAgentChat block): + %s + User request: %s """.formatted( @@ -278,6 +310,7 @@ public class FlowAssistantPromptService { summarizeFlow(currentFlow), summarizeErrors(errors), toJson(descriptor), + buildMcpServerCatalog(mcpServers), userPrompt == null || userPrompt.isBlank() ? "(none)" : userPrompt); } 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 3843235..e554c3e 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 @@ -246,6 +246,9 @@ public class FlowAssistantService { @Autowired private FlowAssistantPromptService promptService; + @Autowired + private it.cnr.isti.workflow.manager.mcp.MCPServersProvider mcpServersProvider; + @Autowired private List> blockFactories; @@ -418,7 +421,8 @@ public class FlowAssistantService { if (targetedRepair) { parsedPlan = buildReusedPlanForTargetedRepair(currentFlow, errors); } else { - String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog); + String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog, + mcpServerCatalogEntries()); parsedPlan = invokeStructuredAndValidate(provider, planningModelFor(mode, phaseModels), phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { ParsedPlan plan = parsePlan(rawResponse); @@ -470,7 +474,7 @@ public class FlowAssistantService { } else { progressListener.onProgress("configuring_blocks", "Configuring block " + blockPlan.blockId()); String blockPrompt = promptService.buildBlockConfigurationPrompt(mode, userPrompt, descriptor, - parsedPlan.plan(), blockPlan, currentFlow, errors, workflowModel); + parsedPlan.plan(), blockPlan, currentFlow, errors, workflowModel, mcpServerCatalogEntries()); int currentBlockIndex = blockIndex; int blockPlanCount = blockPlans.size(); ConfiguredBlockResult configuredBlock = invokeStructuredAndValidate(provider, @@ -665,7 +669,8 @@ public class FlowAssistantService { progressListener.onProgress("configuring_blocks", "Configuring block " + blockPlan.blockId() + " in container " + containerPlan.containerId()); String blockPrompt = promptService.buildBlockConfigurationPrompt(mode, userPrompt, descriptor, - containerInnerPlan(containerPlan), blockPlan, null, List.of(), workflowModel); + containerInnerPlan(containerPlan), blockPlan, null, List.of(), workflowModel, + mcpServerCatalogEntries()); int currentBlockIndex = blockIndex; int blockPlanCount = innerBlockPlans.size(); ConfiguredBlockResult configuredBlock = invokeStructuredAndValidate(provider, @@ -1154,6 +1159,7 @@ public class FlowAssistantService { ensureRequiredTextDefaults(normalizedConfig, descriptor, blockPlan); normalizeHumanDecisionOptions(normalizedConfig, descriptor); normalizeHttpServerCallAuthorization(normalizedConfig, blockPlan); + normalizeMcpAgentServers(normalizedConfig, blockPlan); normalizeMcpAgentSharedMemory(normalizedConfig, blockPlan, model, requireSharedMemorySemantics); ensureSequentialInputPlaceholder(normalizedConfig, descriptor, blockPlan, blockIndex, blockCount); @@ -1224,6 +1230,43 @@ public class FlowAssistantService { || 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 = textOrEmpty(server.path("sourceType")); + boolean catalogSourced = sourceType.isBlank() || "CATALOG".equalsIgnoreCase(sourceType); + if (!catalogSourced) { + kept.add(server); + continue; + } + String serverName = 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())) { @@ -1284,6 +1327,22 @@ public class FlowAssistantService { 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()) { diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java index b9a5ebb..034acb6 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java @@ -781,6 +781,109 @@ public class AssistantControllerTest { assertEquals(2, response.flow().flow().getBlocks().size()); } + @Test + public void draftBindsModelChosenMcpCatalogServerForToolStep() { + // The planner and block-config phases are shown the real MCP catalog, and the model + // binds a concrete catalog server (coding-agent-mcp) to an MCPAgent tool step. That + // explicit choice must be preserved verbatim - never silently replaced by the rag default. + mockCodingAgentToolResponses("coding-agent-mcp"); + + AssistantFlowResponse response = assistantController.draft( + new AssistantGenerationRequest( + "write a program to the workspace and compile it", + MODEL, + 1)); + + assertNotNull(response); + MCPAgentBlockConfiguration configuration = onlyMcpAgentConfiguration(response); + assertNotNull(configuration.getMcpServers()); + assertTrue(configuration.getMcpServers().stream() + .anyMatch(server -> "coding-agent-mcp".equals(server.serverName())), + "expected the model-chosen catalog server to be preserved"); + assertTrue(configuration.getMcpServers().stream() + .noneMatch(server -> "rag".equals(server.serverName())), + "the model's explicit server choice must not be overwritten by the rag default"); + } + + @Test + public void draftDropsUnknownMcpServerChosenByModel() { + // A hallucinated server id must be dropped during normalization so it degrades to a clean + // validation/fallback path instead of a runtime "Unknown MCP server" failure. + mockCodingAgentToolResponses("totally-made-up-server"); + + AssistantFlowResponse response = assistantController.draft( + new AssistantGenerationRequest( + "write a program to the workspace and compile it", + MODEL, + 1)); + + assertNotNull(response); + MCPAgentBlockConfiguration configuration = onlyMcpAgentConfiguration(response); + boolean hasUnknown = configuration.getMcpServers() != null + && configuration.getMcpServers().stream() + .anyMatch(server -> "totally-made-up-server".equals(server.serverName())); + assertFalse(hasUnknown, "an unknown/hallucinated MCP server id must be dropped"); + } + + private MCPAgentBlockConfiguration onlyMcpAgentConfiguration(AssistantFlowResponse response) { + Block agentBlock = response.flow().flow().getBlocks().stream() + .filter(block -> MCPAgentBlockType.TYPE.equals(block.getType().getName())) + .findFirst() + .orElseThrow(() -> new AssertionError("expected an MCPAgent block in the drafted flow")); + return (MCPAgentBlockConfiguration) agentBlock.getSpecificConfiguration(); + } + + private void mockCodingAgentToolResponses(String serverId) { + Answer answer = invocation -> { + String prompt = invocation.getArgument(1, String.class); + if (prompt.contains("TASK: PLAN")) { + // The MCP catalog must be surfaced to the planner so it can reach for a tool step. + assertTrue(prompt.contains("Available MCP servers"), + "plan prompt must list the MCP catalog"); + assertTrue(prompt.contains("coding-agent-mcp"), + "plan prompt must include the coding-agent server from the catalog"); + return TestAssistantResponses.wrap(java.util.Map.of( + "rationale", "One MCP tool step writes code and compiles it.", + "plan", java.util.Map.of( + "name", "Build program", + "description", "Write a program to the workspace and compile it.", + "blocks", java.util.List.of( + java.util.Map.of( + "blockId", "b1", + "blockType", "MCPAgent", + "purpose", "Write the program to the workspace and compile it"))))); + } + if (prompt.contains("TASK: BLOCK_CONFIG")) { + // The catalog must also reach block configuration so the model can bind a server. + assertTrue(prompt.contains("Available MCP servers"), + "block config prompt must list the MCP catalog"); + return TestAssistantResponses.wrap(java.util.Map.of( + "rationale", "Bind the coding-agent MCP server to write files and compile.", + "block", java.util.Map.of( + "blockId", "b1", + "name", "Build program", + "config", java.util.Map.of( + "prompt", "Write a hello-world program to the workspace and compile it.", + "mcpServers", java.util.List.of( + java.util.Map.of( + "sourceType", "CATALOG", + "serverName", serverId, + "configuration", java.util.Map.of())))))); + } + if (prompt.contains("TASK: CONNECTIONS")) { + return TestAssistantResponses.wrap(java.util.Map.of( + "rationale", "Single block; no connections.", + "connections", java.util.List.of())); + } + throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); + }; + + Mockito.when(internalOllamaLLMProvider.generate(Mockito.eq(MODEL), Mockito.anyString())) + .thenAnswer(answer); + Mockito.when(internalOllamaLLMProvider.generateJson(Mockito.eq(MODEL), Mockito.anyString())) + .thenAnswer(answer); + } + private void assertSharedMemoryFlowIsExecutable(FlowCreateRequest flow) { Block indexBlock = flow.flow().getBlocks().stream() .filter(block -> "Index file".equals(block.getName()))