Fix shared MCP sessions in container subflows

This commit is contained in:
Lucio Lelii 2026-09-09 16:34:05 +02:00
parent 289bcfb1db
commit ec009411e1
3 changed files with 93 additions and 6 deletions

View File

@ -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<String> 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<String, Object> projectContext) {
flowExecutionValidator.validate(flow);
Set<String> 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<String, Object> projectContext,
Set<String> externalSessions) {
flowExecutionValidator.validate(flow, externalSessions);
List<ExecutionAuthorizationRequirement> requiredAuthorizations = AuthorizationRequirementResolver.resolveRequiredAuthorizations(flow, llmProviders);
ExecutionObject execObject = ExecutionObject.builder()
.executionName(executionName)

View File

@ -43,7 +43,17 @@ public class FlowExecutionValidator {
Validator validator;
public void validate(FlowData flowData) {
List<ValidationError> 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<String> externalSessions) {
List<ValidationError> errors = collectErrors(flowData,
externalSessions == null ? Set.of() : Set.copyOf(externalSessions));
if (!errors.isEmpty()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, ValidationErrorCodec.encode(errors));

View File

@ -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<GenericContainerType> 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<MCPAgentBlockType> 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<EndBlockType> endBlock(String name) {
return endBlockFactory.create(EndBlockConfiguration.builder()
.name(name)