Fix shared MCP session ownership in subflows
This commit is contained in:
parent
a2743c6c66
commit
760554c2d7
|
|
@ -813,7 +813,7 @@ public class ExecutionsService {
|
|||
private void reconcileGenericSubflow(ExecutionObject parent, Step<?> parentStep,
|
||||
GenericContainerConfiguration configuration, ExecutionObject child) {
|
||||
if (child.getContext().getStatus() == ExecutionStatus.SUCCESS) {
|
||||
parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, child);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs =
|
||||
ContainerFlowInterfaceResolver.getExposedOutputs(configuration.getSubFlow());
|
||||
// routed(), not completed(): a successful subflow that did not take an internal
|
||||
|
|
@ -831,6 +831,12 @@ public class ExecutionsService {
|
|||
+ (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors()));
|
||||
}
|
||||
|
||||
private void mergeSubflowExecutionVariableDescriptors(ExecutionObject parent, ExecutionObject child) {
|
||||
parent.setExecutionVariableDescriptors(SubflowExecutionVariables.mergedBackToParent(
|
||||
parent.getContext().getExecutionVariableDescriptors(),
|
||||
child.getContext().getExecutionVariableDescriptors()));
|
||||
}
|
||||
|
||||
// ---- Iterator: sequential per-item children, no guard ----
|
||||
|
||||
/**
|
||||
|
|
@ -868,7 +874,7 @@ public class ExecutionsService {
|
|||
+ (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors()));
|
||||
return;
|
||||
}
|
||||
parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, child);
|
||||
IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> outputHandles = resolution.resolvedOutputs().stream()
|
||||
.map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList();
|
||||
|
|
@ -945,7 +951,7 @@ public class ExecutionsService {
|
|||
throw new IllegalStateException("IteratorContainer subflow ended with status " + child.getContext().getStatus()
|
||||
+ (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors()));
|
||||
}
|
||||
parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, child);
|
||||
Map<String, Object> iterationOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(child, outputHandles);
|
||||
for (var output : resolution.resolvedOutputs()) {
|
||||
accumulatedOutputs.get(output.publicName()).add(iterationOutputs.get(output.publicName()));
|
||||
|
|
@ -988,7 +994,7 @@ public class ExecutionsService {
|
|||
+ (completedChild.getContext().getErrors().isEmpty() ? "" : ": " + completedChild.getContext().getErrors()));
|
||||
return;
|
||||
}
|
||||
parent.setExecutionVariableDescriptors(completedChild.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, completedChild);
|
||||
|
||||
Map<String, Object> state = continuation.getState() == null ? Map.of() : continuation.getState();
|
||||
Map<String, Object> currentInputs = ContainerExecutionSupport.asStringObjectMap(state.get("currentInputs"));
|
||||
|
|
@ -1081,7 +1087,7 @@ public class ExecutionsService {
|
|||
throw new IllegalStateException("LoopContainer subflow ended with status " + mainChild.getContext().getStatus()
|
||||
+ (mainChild.getContext().getErrors().isEmpty() ? "" : ": " + mainChild.getContext().getErrors()));
|
||||
}
|
||||
parent.setExecutionVariableDescriptors(mainChild.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, mainChild);
|
||||
latestOutputs = ContainerExecutionSupport.collectExposedOutputsAsMap(mainChild, mainOutputHandles);
|
||||
phase = LoopPhase.GUARD;
|
||||
continue;
|
||||
|
|
@ -1103,7 +1109,7 @@ public class ExecutionsService {
|
|||
throw new IllegalStateException("LoopContainer guard subflow ended with status " + guardChild.getContext().getStatus()
|
||||
+ (guardChild.getContext().getErrors().isEmpty() ? "" : ": " + guardChild.getContext().getErrors()));
|
||||
}
|
||||
parent.setExecutionVariableDescriptors(guardChild.getContext().getExecutionVariableDescriptors());
|
||||
mergeSubflowExecutionVariableDescriptors(parent, guardChild);
|
||||
Map<String, Object> guardResult = ContainerExecutionSupport.collectExposedOutputsAsMap(guardChild, guardOutputHandles);
|
||||
boolean shouldContinue = LoopGuardSupport.parseLoopGuardResponse(
|
||||
String.valueOf(LoopGuardSupport.requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT)));
|
||||
|
|
@ -1152,7 +1158,7 @@ public class ExecutionsService {
|
|||
propagateAuthorizationValue(child.getId(), key, authorizations.get(key));
|
||||
}
|
||||
}
|
||||
setExecutionVariableDescriptors(child.getId(), parent.getContext().getExecutionVariableDescriptors());
|
||||
setSubflowInheritedExecutionVariableDescriptors(child.getId(), parent.getContext().getExecutionVariableDescriptors());
|
||||
setGlobalInputDescriptors(child.getId(),
|
||||
SubflowGlobalInputs.descriptorsFor(child, parent.getContext().getGlobalInputs()));
|
||||
for (var entry : inputPortsByName.entrySet()) {
|
||||
|
|
@ -1341,6 +1347,13 @@ public class ExecutionsService {
|
|||
return eo;
|
||||
}
|
||||
|
||||
/** Copies variables into a container child without transferring ownership of MCP sessions. */
|
||||
public ExecutionObject setSubflowInheritedExecutionVariableDescriptors(String executionId,
|
||||
Map<String, ExecutionVariableDescriptor> parentExecutionVariableDescriptors) {
|
||||
return setExecutionVariableDescriptors(executionId,
|
||||
SubflowExecutionVariables.inheritedByChild(parentExecutionVariableDescriptors));
|
||||
}
|
||||
|
||||
public ExecutionObject setGlobalInputDescriptors(String executionId,
|
||||
Map<String, ExecutionVariableDescriptor> globalInputDescriptors) {
|
||||
ExecutionObject eo = getExecution(executionId);
|
||||
|
|
|
|||
|
|
@ -0,0 +1,67 @@
|
|||
package it.cnr.isti.workflow.manager.executions;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
|
||||
/** Ownership rules for execution variables crossing a container boundary. */
|
||||
public final class SubflowExecutionVariables {
|
||||
|
||||
private SubflowExecutionVariables() {
|
||||
}
|
||||
|
||||
/**
|
||||
* A child may use a parent-owned MCP session but must never close it when the child ends.
|
||||
*/
|
||||
public static Map<String, ExecutionVariableDescriptor> inheritedByChild(
|
||||
Map<String, ExecutionVariableDescriptor> parentDescriptors) {
|
||||
LinkedHashMap<String, ExecutionVariableDescriptor> inherited = new LinkedHashMap<>();
|
||||
if (parentDescriptors == null) {
|
||||
return inherited;
|
||||
}
|
||||
parentDescriptors.forEach((key, descriptor) -> {
|
||||
if (descriptor == null) {
|
||||
return;
|
||||
}
|
||||
ExecutionVariableDescriptor normalized = ExecutionVariableRegistry.normalize(descriptor);
|
||||
if (normalized.getKind() == ExecutionVariableKind.MCP_SESSION) {
|
||||
normalized.setCleanupPolicy(ExecutionVariableCleanupPolicy.NONE);
|
||||
}
|
||||
inherited.put(normalized.getName(), normalized);
|
||||
});
|
||||
return inherited;
|
||||
}
|
||||
|
||||
/**
|
||||
* Carries a completed child's variables back to its parent. An MCP session which the parent
|
||||
* already owned retains the parent's close-on-final-execution policy.
|
||||
*/
|
||||
public static Map<String, ExecutionVariableDescriptor> mergedBackToParent(
|
||||
Map<String, ExecutionVariableDescriptor> parentDescriptors,
|
||||
Map<String, ExecutionVariableDescriptor> childDescriptors) {
|
||||
LinkedHashMap<String, ExecutionVariableDescriptor> merged = new LinkedHashMap<>();
|
||||
if (childDescriptors != null) {
|
||||
childDescriptors.values().stream()
|
||||
.filter(Objects::nonNull)
|
||||
.map(ExecutionVariableRegistry::normalize)
|
||||
.forEach(descriptor -> merged.put(descriptor.getName(), descriptor));
|
||||
}
|
||||
if (parentDescriptors == null) {
|
||||
return merged;
|
||||
}
|
||||
parentDescriptors.values().stream()
|
||||
.filter(Objects::nonNull)
|
||||
.map(ExecutionVariableRegistry::normalize)
|
||||
.filter(descriptor -> descriptor.getKind() == ExecutionVariableKind.MCP_SESSION
|
||||
&& descriptor.getCleanupPolicy() == ExecutionVariableCleanupPolicy.CLOSE_RESOURCE)
|
||||
.forEach(parentDescriptor -> {
|
||||
ExecutionVariableDescriptor childDescriptor = merged.get(parentDescriptor.getName());
|
||||
if (childDescriptor != null
|
||||
&& childDescriptor.getKind() == ExecutionVariableKind.MCP_SESSION
|
||||
&& Objects.equals(childDescriptor.getValue(), parentDescriptor.getValue())) {
|
||||
merged.put(parentDescriptor.getName(), parentDescriptor);
|
||||
}
|
||||
});
|
||||
return merged;
|
||||
}
|
||||
}
|
||||
|
|
@ -25,6 +25,7 @@ import it.cnr.isti.workflow.manager.executions.FieldKey;
|
|||
import it.cnr.isti.workflow.manager.executions.ContainerSubflowRole;
|
||||
import it.cnr.isti.workflow.manager.executions.NodeExecutionResult;
|
||||
import it.cnr.isti.workflow.manager.executions.SubflowGlobalInputs;
|
||||
import it.cnr.isti.workflow.manager.executions.SubflowExecutionVariables;
|
||||
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
|
||||
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasContainerPropagation;
|
||||
import it.cnr.isti.workflow.manager.executions.persistence.ContainerContinuationPhase;
|
||||
|
|
@ -70,7 +71,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
.build());
|
||||
|
||||
propagateAuthorizations(innerExecution, authorizations);
|
||||
executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors);
|
||||
executionsService.setSubflowInheritedExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors);
|
||||
executionsService.setGlobalInputDescriptors(innerExecution.getId(), SubflowGlobalInputs.descriptorsFor(
|
||||
innerExecution, ExecutionRuntimeContextSupport.globalView(executionVariables)));
|
||||
|
||||
|
|
@ -121,8 +122,10 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
}
|
||||
executionVariables.clear();
|
||||
executionVariables.putAll(innerExecution.getContext().getExecutionVariables());
|
||||
Map<String, ExecutionVariableDescriptor> mergedDescriptors = SubflowExecutionVariables.mergedBackToParent(
|
||||
executionVariableDescriptors, innerExecution.getContext().getExecutionVariableDescriptors());
|
||||
executionVariableDescriptors.clear();
|
||||
executionVariableDescriptors.putAll(innerExecution.getContext().getExecutionVariableDescriptors());
|
||||
executionVariableDescriptors.putAll(mergedDescriptors);
|
||||
|
||||
Map<String, Object> outputs = new java.util.LinkedHashMap<>();
|
||||
java.util.Set<String> exposedNames = new java.util.LinkedHashSet<>();
|
||||
|
|
|
|||
|
|
@ -0,0 +1,30 @@
|
|||
package it.cnr.isti.workflow.manager.executions;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
class SubflowExecutionVariablesTest {
|
||||
|
||||
@Test
|
||||
void childCannotCloseSessionInheritedFromParentAndParentRetainsOwnershipOnMerge() {
|
||||
ExecutionVariableDescriptor parentSession = ExecutionVariableDescriptor.builder()
|
||||
.name("mcp-coding")
|
||||
.kind(ExecutionVariableKind.MCP_SESSION)
|
||||
.value("session-1")
|
||||
.cleanupPolicy(ExecutionVariableCleanupPolicy.CLOSE_RESOURCE)
|
||||
.build();
|
||||
|
||||
Map<String, ExecutionVariableDescriptor> childDescriptors = SubflowExecutionVariables.inheritedByChild(
|
||||
Map.of("mcp-coding", parentSession));
|
||||
assertEquals(ExecutionVariableCleanupPolicy.NONE,
|
||||
childDescriptors.get("mcp-coding").getCleanupPolicy());
|
||||
|
||||
Map<String, ExecutionVariableDescriptor> parentAfterChild = SubflowExecutionVariables.mergedBackToParent(
|
||||
Map.of("mcp-coding", parentSession), childDescriptors);
|
||||
assertEquals(ExecutionVariableCleanupPolicy.CLOSE_RESOURCE,
|
||||
parentAfterChild.get("mcp-coding").getCleanupPolicy());
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue