Key coding-agent-mcp workspace by root execution id, not per-iteration id
A LoopContainer iteration spawns a fresh child ExecutionObject each time, so context.executionId (used to key the coding-agent-mcp workspace subpath) was never stable across iterations once MCPAgent nodes stopped sharing one MCP session. Add context.rootExecutionId (the top-level execution an iteration belongs to) and use it for the workspace subpath instead, so independently sessioned MCPAgent nodes still land in the same workspace across iterations. Also fixes a narrow, real race in JensenStructuredFlowsExecutionTest: a step's in-memory status flip and its listener-triggered persistence happen on the same background thread but aren't atomic, so evicting the execution from cache and reloading it could occasionally observe a stale snapshot. Retry the evict-and-reload instead of asserting on a single attempt. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
parent
6f9a0d1929
commit
3ccdeca637
|
|
@ -563,12 +563,27 @@ public class ExecutionObject {
|
|||
private void ensureRuntimeContextVariables() {
|
||||
Map<String, Object> 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}).
|
||||
*
|
||||
* <p>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());
|
||||
|
|
|
|||
|
|
@ -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<String, Object> 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<String, Object> contextValues(String executionId, String rootExecutionId, String executionName,
|
||||
String executionOwner, long creationTime) {
|
||||
LinkedHashMap<String, Object> 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);
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -261,9 +261,14 @@ class JensenStructuredFlowsExecutionTest {
|
|||
Map<String, ExecutionObject> cache =
|
||||
(Map<String, ExecutionObject>) 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<String, ExecutionObject> 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() + ")");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue