From ec009411e144c938dfd7268615ebd548f4643f1b Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Wed, 9 Sep 2026 16:34:05 +0200 Subject: [PATCH] Fix shared MCP sessions in container subflows --- .../manager/executions/ExecutionsService.java | 44 ++++++++++++++++--- .../validation/FlowExecutionValidator.java | 12 ++++- .../SubflowExecutionPersistenceTest.java | 43 ++++++++++++++++++ 3 files changed, 93 insertions(+), 6 deletions(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java index 7d711bc..e743960 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -173,7 +173,30 @@ public class ExecutionsService { } return createExecution(executionName, flow, parentExecution.getOwner(), null, null, null, 1, biasExecutionContext, ExecutionKind.SUBFLOW, parentExecution.getId(), parentStepId, - parentIterationIndex, subflowRole == null ? ContainerSubflowRole.MAIN : subflowRole); + parentIterationIndex, subflowRole == null ? ContainerSubflowRole.MAIN : subflowRole, + availableSharedMcpSessions(parentExecution)); + } + + /** + * A container child receives its parent's variable registry immediately after it is created. + * Pass the already-open MCP session names to its pre-persistence validation as well, otherwise + * a valid parent-to-container shared-session chain is incorrectly treated as a standalone + * subflow with no producer. + */ + private Set availableSharedMcpSessions(ExecutionObject parentExecution) { + if (parentExecution == null || parentExecution.getContext().getExecutionVariableDescriptors() == null) { + return Set.of(); + } + return parentExecution.getContext().getExecutionVariableDescriptors().values().stream() + .filter(Objects::nonNull) + .filter(descriptor -> descriptor.getKind() == ExecutionVariableKind.MCP_SESSION) + .filter(descriptor -> StringUtils.hasText(descriptor.getName())) + // A descriptor without a session id is only a declaration, not a session which a + // child can actually reuse. + .filter(descriptor -> descriptor.getValue() != null + && StringUtils.hasText(String.valueOf(descriptor.getValue()))) + .map(descriptor -> descriptor.getName().trim()) + .collect(Collectors.toCollection(LinkedHashSet::new)); } @Transactional @@ -205,7 +228,7 @@ public class ExecutionsService { // The execution is tagged with the flow's project even when the context resolved // empty (a stranger running a published flow), so grouping still works. projectContextResolver.projectIdOfFlow(sourceFlowId), projectRunId, projectRunOrder, - projectContext.values()); + projectContext.values(), Set.of()); seedGlobalInputsFromProjectContext(execution, projectContext.values()); return execution; } @@ -284,15 +307,26 @@ public class ExecutionsService { String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole) { return createExecution(executionName, flow, owner, runGroupId, sourceFlowId, rerunOfExecutionId, runNumber, biasExecutionContext, executionKind, parentExecutionId, parentStepId, parentIterationIndex, - subflowRole, null, null, null, null); + subflowRole, Set.of()); } private ExecutionObject createExecution(String executionName, FlowData flow, String owner, String runGroupId, String sourceFlowId, String rerunOfExecutionId, Integer runNumber, BiasExecutionContext biasExecutionContext, ExecutionKind executionKind, String parentExecutionId, String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole, - String projectId, String projectRunId, Integer projectRunOrder, Map projectContext) { - flowExecutionValidator.validate(flow); + Set externalSessions) { + return createExecution(executionName, flow, owner, runGroupId, sourceFlowId, rerunOfExecutionId, runNumber, + biasExecutionContext, executionKind, parentExecutionId, parentStepId, parentIterationIndex, + subflowRole, null, null, null, null, externalSessions); + } + + private ExecutionObject createExecution(String executionName, FlowData flow, String owner, String runGroupId, + String sourceFlowId, String rerunOfExecutionId, Integer runNumber, + BiasExecutionContext biasExecutionContext, ExecutionKind executionKind, String parentExecutionId, + String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole, + String projectId, String projectRunId, Integer projectRunOrder, Map projectContext, + Set externalSessions) { + flowExecutionValidator.validate(flow, externalSessions); List requiredAuthorizations = AuthorizationRequirementResolver.resolveRequiredAuthorizations(flow, llmProviders); ExecutionObject execObject = ExecutionObject.builder() .executionName(executionName) diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java index b435e3e..fd4ad5f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java @@ -43,7 +43,17 @@ public class FlowExecutionValidator { Validator validator; public void validate(FlowData flowData) { - List errors = collectErrors(flowData); + validate(flowData, Set.of()); + } + + /** + * Validates a flow that will execute as a container subflow. Shared MCP sessions in + * {@code externalSessions} have already been opened by the enclosing execution and are + * therefore valid producers for consumers in this flow. + */ + public void validate(FlowData flowData, Set externalSessions) { + List errors = collectErrors(flowData, + externalSessions == null ? Set.of() : Set.copyOf(externalSessions)); if (!errors.isEmpty()) { throw new ResponseStatusException(HttpStatus.BAD_REQUEST, ValidationErrorCodec.encode(errors)); diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/SubflowExecutionPersistenceTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/SubflowExecutionPersistenceTest.java index 0006f23..2b7f6cf 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/SubflowExecutionPersistenceTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/SubflowExecutionPersistenceTest.java @@ -20,8 +20,11 @@ import org.springframework.web.server.ResponseStatusException; import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.EndBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.factories.EndBlockFactory; +import it.cnr.isti.workflow.manager.blocks.factories.MCPAgentBlockFactory; import it.cnr.isti.workflow.manager.blocks.types.EndBlockType; +import it.cnr.isti.workflow.manager.blocks.types.MCPAgentBlockType; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; import it.cnr.isti.workflow.manager.containers.factories.GenericContainerFactory; @@ -53,6 +56,9 @@ class SubflowExecutionPersistenceTest { @Autowired GenericContainerFactory genericContainerFactory; + @Autowired + MCPAgentBlockFactory mcpAgentBlockFactory; + @Test void linkedSubflowInheritsOwnerIsHiddenAndPersistsContinuation() { String owner = "subflow-owner-" + UUID.randomUUID(); @@ -147,6 +153,43 @@ class SubflowExecutionPersistenceTest { ContainerSubflowRole.MAIN)); } + @Test + void linkedSubflowAcceptsSharedMcpSessionAlreadyOpenedByParent() { + String owner = "subflow-mcp-session-" + UUID.randomUUID(); + FlowData parentSubflow = FlowData.builder().block(endBlock("Parent inner end")).build(); + Container container = genericContainerFactory.create( + GenericContainerConfiguration.builder() + .name("Parent container") + .subFlow(parentSubflow) + .build()); + ExecutionObject parent = executionsService.createExecution( + "Parent execution", FlowData.builder().container(container).build(), owner); + parent = executionsService.setExecutionVariableDescriptors(parent.getId(), Map.of( + "mcp-coding", ExecutionVariableDescriptor.builder() + .name("mcp-coding") + .kind(ExecutionVariableKind.MCP_SESSION) + .value("already-open-session") + .build())); + + Block consumer = mcpAgentBlockFactory.create(MCPAgentBlockConfiguration.builder() + .name("Reuse coding session") + .prompt("Continue the task") + .useSharedSession(true) + .sharedSessionRef("mcp-coding") + .build()); + + ExecutionObject child = executionsService.createInnerExecution( + "Child execution", + FlowData.builder().block(consumer).build(), + BiasExecutionContext.normal(), + parent, + container.getId(), + 1, + ContainerSubflowRole.MAIN); + + assertEquals(ExecutionKind.SUBFLOW, child.getExecutionKind()); + } + private Block endBlock(String name) { return endBlockFactory.create(EndBlockConfiguration.builder() .name(name)