Close shared MCP sessions when their execution is cancelled
ExecutionContext.cancel() cleared executionVariableDescriptors before
transitioning to CANCELLED, so ExecutionsService.cleanupManagedResourcesIfFinal
(triggered from the same state-change notification) always found an empty
map and closed nothing. The CLOSE_RESOURCE mechanism already worked
correctly on SUCCESS and ERROR - only cancel/stop silently leaked any
shared MCP bridge session still open at that point, until the bridge's
own idle timeout, eventually hitting its session cap ("Reached the
maximum limit of 5 sessions").
Defer clearing executionVariableDescriptors until after the state-change
notification runs, so the cleanup sees the still-registered CLOSE_RESOURCE
session descriptor and closes it before the map is scrubbed. Added a
regression test that reproduces the leak (confirmed red without this
change) and asserts closeSessionQuietly runs on cancel.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
parent
44991f4a79
commit
cc820554ac
|
|
@ -555,7 +555,6 @@ public class ExecutionContext implements ExecutionListener {
|
|||
this.result.clear();
|
||||
this.partialResult.clear();
|
||||
this.executionVariables.clear();
|
||||
this.executionVariableDescriptors.clear();
|
||||
this.globalInputs.clear();
|
||||
this.globalInputDescriptors.clear();
|
||||
this.projectContext.clear();
|
||||
|
|
@ -568,7 +567,14 @@ public class ExecutionContext implements ExecutionListener {
|
|||
this.waitingSteps.clear();
|
||||
this.steps.values().forEach(Step::cancel);
|
||||
this.setStatus(ExecutionStatus.CANCELLED);
|
||||
// executionVariableDescriptors is cleared only after this notification, not before: the
|
||||
// stateChangeListener synchronously runs ExecutionsService.cleanupManagedResourcesIfFinal,
|
||||
// which reads this exact map to find and close any CLOSE_RESOURCE-tagged MCP session
|
||||
// (see MCPSharedSessionRegistry). Clearing it first - as every other field here does -
|
||||
// left that lookup with nothing to find, so a cancelled execution never closed its shared
|
||||
// MCP bridge sessions and they leaked until the bridge's own idle timeout.
|
||||
notifyStateChanged();
|
||||
this.executionVariableDescriptors.clear();
|
||||
}
|
||||
|
||||
public void setStateChangeListener(Runnable stateChangeListener) {
|
||||
|
|
|
|||
|
|
@ -1305,6 +1305,64 @@ public class ExecutionTest {
|
|||
Mockito.verify(mcpAgentService, Mockito.timeout(1000)).closeSessionQuietly("shared-session-1");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void cancellingAnExecutionClosesASharedMcpSessionOpenedEarlierInIt() {
|
||||
// ExecutionContext.cancel() used to clear executionVariableDescriptors before notifying
|
||||
// listeners, so the CLOSE_RESOURCE cleanup in ExecutionsService.cleanupManagedResourcesIfFinal
|
||||
// always found an empty map and closed nothing on cancel - only on SUCCESS/ERROR. A shared MCP
|
||||
// session opened by an earlier block (like "Research session" below) and still referenced by a
|
||||
// later, still-waiting block would then leak until the bridge's own idle timeout.
|
||||
Block<MCPAgentBlockType> producer = mcpAgentBlockFactory.create(MCPAgentBlockConfiguration.builder()
|
||||
.name("Research session")
|
||||
.model("llama3.1:8b")
|
||||
.prompt("Find data for ${{cand}}")
|
||||
.shareSession(true)
|
||||
.sharedSessionName("candidateResearchSession")
|
||||
.build());
|
||||
|
||||
Block<MCPAgentChatBlockType> chatBlock = mcpAgentChatBlockFactory.create(MCPAgentChatBlockConfiguration.builder()
|
||||
.name("MCP Chat")
|
||||
.model("llama3.1:8b")
|
||||
.goalDescription("Assess ${{cand}} and reach a final decision")
|
||||
.inputs(List.of(new ChatInteractionInput("cand", IOType.TEXT, false)))
|
||||
.useSharedSession(true)
|
||||
.sharedSessionRef("candidateResearchSession")
|
||||
.build());
|
||||
|
||||
Flow flow = Flow.builder()
|
||||
.name("Shared MCP session cancelled mid-flight")
|
||||
.description("A later block still waits on the session an earlier block opened and shared")
|
||||
.block(producer)
|
||||
.block(chatBlock)
|
||||
.connection(Connection.builder()
|
||||
.sourceId(producer.getId())
|
||||
.sourceName(MCPAgentBlockFactory.OUTPUT_NAME)
|
||||
.targetId(chatBlock.getId())
|
||||
.targetName("cand")
|
||||
.build())
|
||||
.build();
|
||||
|
||||
Mockito.when(mcpAgentService.openSession(Mockito.eq("llama3.1:8b"), Mockito.anyList(), Mockito.anyMap(), Mockito.any()))
|
||||
.thenReturn("shared-session-cancel");
|
||||
Mockito.when(mcpAgentService.querySession("shared-session-cancel", "Find data for Ada"))
|
||||
.thenReturn("Ada");
|
||||
|
||||
ExecutionObject execObject = executionsService.createExecution(flow);
|
||||
execObject = executionsService.prepareInput(execObject.getId(), producer.getId(), "cand", "Ada");
|
||||
execObject = executionsService.startExecution(execObject.getId());
|
||||
while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) {
|
||||
execObject = executionsService.getExecution(execObject.getId());
|
||||
}
|
||||
|
||||
assertEquals(ExecutionStatus.WAITING, execObject.getContext().getStatus());
|
||||
assertEquals("shared-session-cancel", execObject.getContext().getExecutionVariables().get("candidateResearchSession"));
|
||||
|
||||
execObject = executionsService.cancelExecution(execObject.getId());
|
||||
|
||||
assertEquals(ExecutionStatus.CANCELLED, execObject.getContext().getStatus());
|
||||
Mockito.verify(mcpAgentService, Mockito.timeout(1000)).closeSessionQuietly("shared-session-cancel");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void failingStepAbortsOtherRunningOrWaitingBranches() {
|
||||
Block<MCPAgentBlockType> failingBlock = mcpAgentBlockFactory.create(MCPAgentBlockConfiguration.builder()
|
||||
|
|
|
|||
Loading…
Reference in New Issue