From cc820554ac1ab4ba2fe0d5c2ee0335e3b362f1c5 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 14 Sep 2026 11:39:01 +0200 Subject: [PATCH] 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 --- .../manager/executions/ExecutionContext.java | 8 ++- .../manager/executions/ExecutionTest.java | 58 +++++++++++++++++++ 2 files changed, 65 insertions(+), 1 deletion(-) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java index 93bc5b8..2f457bf 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java @@ -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) { diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java index 77e55d4..e4a80d4 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java @@ -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 producer = mcpAgentBlockFactory.create(MCPAgentBlockConfiguration.builder() + .name("Research session") + .model("llama3.1:8b") + .prompt("Find data for ${{cand}}") + .shareSession(true) + .sharedSessionName("candidateResearchSession") + .build()); + + Block 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 failingBlock = mcpAgentBlockFactory.create(MCPAgentBlockConfiguration.builder()