diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java index 3c4c48c..e12d8b1 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java @@ -563,12 +563,27 @@ public class ExecutionObject { private void ensureRuntimeContextVariables() { Map runtimeContext = ExecutionRuntimeContextSupport.contextValues( this.id, + getRootExecutionId(), this.name, this.owner, this.creationTime); this.context.setRuntimeContextValues(runtimeContext); } + /** + * The top-level execution this run belongs to: itself for a {@link ExecutionKind#TOP_LEVEL} + * execution, or its immediate parent for a container subflow (subflow executions are always + * direct children of a top-level execution; nesting deeper than one level is not supported, see + * {@link ExecutionsService#createInnerExecution}). + * + *

Exposed as {@code context.rootExecutionId} so an MCP session opened without + * {@code useSharedSession} can still be pinned to the same workspace as the top-level run, + * instead of the fresh execution id each loop iteration gets. + */ + public String getRootExecutionId() { + return this.parentExecutionId != null ? this.parentExecutionId : this.id; + } + private void registerGlobalInputs() { for (IODescriptor globalInput : this.requiredGlobalInputs) { ExecutionVariableDescriptor existing = this.context.getGlobalInputDescriptor(globalInput.getName()); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java index 34d1389..88602b4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java @@ -11,6 +11,7 @@ public final class ExecutionRuntimeContextSupport { public static final String GLOBAL_PREFIX = "global."; public static final String PROJECT_PREFIX = "project."; public static final String EXECUTION_ID = CONTEXT_PREFIX + "executionId"; + public static final String ROOT_EXECUTION_ID = CONTEXT_PREFIX + "rootExecutionId"; public static final String EXECUTION_NAME = CONTEXT_PREFIX + "executionName"; public static final String EXECUTION_OWNER = CONTEXT_PREFIX + "executionOwner"; public static final String EXECUTION_CREATION_TIME = CONTEXT_PREFIX + "creationTime"; @@ -18,10 +19,18 @@ public final class ExecutionRuntimeContextSupport { private ExecutionRuntimeContextSupport() { } - public static Map contextValues(String executionId, String executionName, + /** + * @param rootExecutionId the top-level execution this run belongs to (itself for a top-level + * execution, its parent for a container subflow) — exposed as + * {@code context.rootExecutionId} so state that must survive across a + * loop's per-iteration executions (e.g. an MCP session's workspace) can + * be keyed on it instead of the current, per-iteration execution id. + */ + public static Map contextValues(String executionId, String rootExecutionId, String executionName, String executionOwner, long creationTime) { LinkedHashMap context = new LinkedHashMap<>(); context.put(EXECUTION_ID, executionId); + context.put(ROOT_EXECUTION_ID, rootExecutionId); context.put(EXECUTION_NAME, executionName); context.put(EXECUTION_OWNER, executionOwner); context.put(EXECUTION_CREATION_TIME, creationTime); diff --git a/src/main/resources/mcp-servers.json b/src/main/resources/mcp-servers.json index dcce7f5..c9cea58 100644 --- a/src/main/resources/mcp-servers.json +++ b/src/main/resources/mcp-servers.json @@ -92,7 +92,7 @@ "sseUrl": "${{host}}/coding-agent/events", "headers": { "Authorization": "Bearer ${{key}}", - "x-workspace-subpath": "${{context.executionId}}" + "x-workspace-subpath": "${{context.rootExecutionId}}" }, "configurationSchema": { "type": "object", diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/JensenStructuredFlowsExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/JensenStructuredFlowsExecutionTest.java index 8059061..a4de994 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/JensenStructuredFlowsExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/JensenStructuredFlowsExecutionTest.java @@ -261,9 +261,14 @@ class JensenStructuredFlowsExecutionTest { Map cache = (Map) ReflectionTestUtils.getField(executionsService, "executions"); assertNotNull(cache); - cache.remove(execution.getId()); - ExecutionObject restored = executionsService.getExecution(execution.getId()); + // Persistence is triggered by a listener the step's background thread calls right after it + // flips its in-memory status to WAITING_FOR_INTERACTION (see Step#run): there is a short, + // genuine window where the in-memory status (what awaitStepStatus above just observed) is + // ahead of the durable snapshot. Evicting from cache and reloading from persistence can land + // inside that window, so retry the evict-and-reload instead of asserting on a single attempt. + ExecutionObject restored = awaitPersistedStepStatus(cache, execution.getId(), exceptionalDecision.getId(), + StepStatus.WAITING_FOR_INTERACTION); assertEquals(ExecutionStatus.SUSPENDED, restored.getContext().getStatus()); assertEquals(StepStatus.WAITING_FOR_INTERACTION, restored.getContext().getSteps().get(exceptionalDecision.getId()).getStatus()); @@ -379,4 +384,27 @@ class JensenStructuredFlowsExecutionTest { } while (Instant.now().isBefore(deadline)); throw new AssertionError("Step " + blockId + " did not reach " + expected); } + + /** + * Evicts the execution from the in-memory cache and reloads it from persistence, retrying for a + * bounded window: the reload can otherwise land in the brief gap between a step's in-memory + * status flip and the listener-triggered persist that follows it (see the comment at the call + * site), landing on a stale snapshot exactly once in a while under load. + */ + private ExecutionObject awaitPersistedStepStatus(Map cache, String executionId, + String blockId, StepStatus expected) { + Instant deadline = Instant.now().plus(Duration.ofSeconds(15)); + ExecutionObject reloaded; + do { + cache.remove(executionId); + reloaded = executionsService.getExecution(executionId); + var step = reloaded.getContext().getSteps().get(blockId); + if (step != null && step.getStatus() == expected) { + return reloaded; + } + Thread.onSpinWait(); + } while (Instant.now().isBefore(deadline)); + throw new AssertionError("Reloaded step " + blockId + " did not reach " + expected + + " (was " + reloaded.getContext().getSteps().get(blockId).getStatus() + ")"); + } }