Say why a subflow cannot start, and propagate its globals from one place
Two follow-ups to the iterator/loop global input fix. **One propagation.** Every container path carried the parent's global inputs into its child with its own copy of the same code, and the copies drifted: the one in ExecutionsService read the unprefixed variables map through a "global." view and handed every iterator and loop child nulls, while the copy in GenericContainerExecutor read the runtime map and kept working. Both now call SubflowGlobalInputs.descriptorsFor, so there is nothing left to diverge. It also sets `multiple` on the descriptor, which neither copy did. The IODescriptor-to-kind switch had grown four identical copies for the same reason, one per descriptor-building site. It now lives on ExecutionVariableKind as forDescriptor. **A message that says something.** Starting a non-READY execution reported only "is not in READY status (CURRENT STATUS is CREATED)". For a container subflow that named an execution the user never sees and gave nothing to act on. What holds an execution back is one of three things, so notStartableReason names it: missing global inputs, missing credentials, or inputs still to provide and on which step. startContainerChild is the single point every container path starts its child from, so it frames that with the container's own name: The subflow of container 'Candidate loop' cannot start: no value for the global input 'who' and the error is filed against the container step, so the diagram points at it. Also fixed: BranchRejoinConcurrencyTest did not compile on its own (a capture conversion the Eclipse compiler rejects), which the full suite had been hiding through incremental compilation. 516 tests green. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
4fa87ae3a2
commit
2ea8e33987
|
|
@ -8,7 +8,6 @@ import java.util.Set;
|
|||
import java.util.stream.Collectors;
|
||||
|
||||
import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver;
|
||||
import it.cnr.isti.workflow.manager.ios.IODescriptor;
|
||||
|
||||
final class ContainerExecutionSupport {
|
||||
|
||||
|
|
@ -35,16 +34,6 @@ final class ContainerExecutionSupport {
|
|||
.collect(Collectors.toCollection(LinkedHashSet::new));
|
||||
}
|
||||
|
||||
static ExecutionVariableKind mapExecutionVariableKind(IODescriptor descriptor) {
|
||||
return switch (descriptor.getType()) {
|
||||
case TEXT -> ExecutionVariableKind.TEXT;
|
||||
case BOOLEAN -> ExecutionVariableKind.BOOLEAN;
|
||||
case FILE, CSV -> ExecutionVariableKind.FILE_PATH;
|
||||
case JSON -> ExecutionVariableKind.JSON;
|
||||
case ANY -> ExecutionVariableKind.ANY;
|
||||
};
|
||||
}
|
||||
|
||||
static Map<String, Object> iteratorState(List<Object> remainingValues, Map<String, Object> runtimeInputValues) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("remainingValues", new java.util.ArrayList<>(remainingValues));
|
||||
|
|
|
|||
|
|
@ -315,9 +315,51 @@ public class ExecutionObject {
|
|||
}
|
||||
if (this.context.getStatus() == ExecutionStatus.READY){
|
||||
this.context.start(executorService);
|
||||
} else
|
||||
throw new IllegalStateException("Execution with id " + this.getId()
|
||||
+ " is not in READY status (CURRENT STATUS is " + this.getContext().getStatus() + ")");
|
||||
} else
|
||||
throw new IllegalStateException("Execution with id " + this.getId() + " cannot start: "
|
||||
+ notStartableReason());
|
||||
}
|
||||
|
||||
/**
|
||||
* Why this execution cannot start, in the user's terms, or null when it can.
|
||||
*
|
||||
* <p>The old message said only "is not in READY status (CURRENT STATUS is CREATED)", which for
|
||||
* a container subflow named an execution the user never sees and gave no way to act. What
|
||||
* actually holds an execution back is one of three things, so it says which.
|
||||
*/
|
||||
public String notStartableReason() {
|
||||
if (this.context.getStatus() == ExecutionStatus.READY) {
|
||||
return null;
|
||||
}
|
||||
List<String> reasons = new ArrayList<>();
|
||||
List<String> missingGlobalInputs = getMissingGlobalInputKeys();
|
||||
if (!missingGlobalInputs.isEmpty()) {
|
||||
reasons.add("no value for the global input" + (missingGlobalInputs.size() == 1 ? " " : "s ")
|
||||
+ quoted(missingGlobalInputs));
|
||||
}
|
||||
List<String> missingAuthorizations = getMissingAuthorizationKeys();
|
||||
if (!missingAuthorizations.isEmpty()) {
|
||||
reasons.add("no credential for " + quoted(missingAuthorizations));
|
||||
}
|
||||
List<String> unsetInputs = new ArrayList<>();
|
||||
for (Step<?> step : this.context.getSteps().values()) {
|
||||
List<String> names = step.getInputs().stream()
|
||||
.filter(input -> !input.isRegistered() && !input.isSet())
|
||||
.map(input -> input.getDescriptor().getName())
|
||||
.toList();
|
||||
if (!names.isEmpty()) {
|
||||
unsetInputs.add(step.getNode().getName() + " (" + String.join(", ", names) + ")");
|
||||
}
|
||||
}
|
||||
if (!unsetInputs.isEmpty()) {
|
||||
reasons.add("inputs still to provide on " + String.join("; ", unsetInputs));
|
||||
}
|
||||
// Nothing above explains it, so the status is all there is to report - and worth knowing.
|
||||
return reasons.isEmpty() ? "its status is " + this.context.getStatus() : String.join("; ", reasons);
|
||||
}
|
||||
|
||||
private static String quoted(List<String> names) {
|
||||
return names.stream().map(name -> "'" + name + "'").collect(java.util.stream.Collectors.joining(", "));
|
||||
}
|
||||
|
||||
protected ExecutionStatus resume() {
|
||||
|
|
@ -533,7 +575,7 @@ public class ExecutionObject {
|
|||
this.context.registerGlobalInput(ExecutionVariableDescriptor.builder()
|
||||
.name(globalInput.getName())
|
||||
.value(existing == null ? null : existing.getValue())
|
||||
.kind(toExecutionVariableKind(globalInput))
|
||||
.kind(ExecutionVariableKind.forDescriptor(globalInput))
|
||||
.multiple(globalInput.isMultiple())
|
||||
.description("Global flow input")
|
||||
.cleanupPolicy(ExecutionVariableCleanupPolicy.NONE)
|
||||
|
|
@ -541,14 +583,4 @@ public class ExecutionObject {
|
|||
}
|
||||
}
|
||||
|
||||
private ExecutionVariableKind toExecutionVariableKind(IODescriptor descriptor) {
|
||||
return switch (descriptor.getType()) {
|
||||
case TEXT -> ExecutionVariableKind.TEXT;
|
||||
case BOOLEAN -> ExecutionVariableKind.BOOLEAN;
|
||||
case FILE, CSV -> ExecutionVariableKind.FILE_PATH;
|
||||
case JSON -> ExecutionVariableKind.JSON;
|
||||
case ANY -> ExecutionVariableKind.ANY;
|
||||
};
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
package it.cnr.isti.workflow.manager.executions;
|
||||
|
||||
import it.cnr.isti.workflow.manager.ios.IODescriptor;
|
||||
|
||||
public enum ExecutionVariableKind {
|
||||
ANY,
|
||||
TEXT,
|
||||
|
|
@ -7,5 +9,19 @@ public enum ExecutionVariableKind {
|
|||
JSON,
|
||||
FILE_PATH,
|
||||
HTTP_RESOURCE,
|
||||
MCP_SESSION
|
||||
MCP_SESSION;
|
||||
|
||||
/**
|
||||
* The kind a declared input maps to. It had grown four identical copies - one per place that
|
||||
* built a variable descriptor from an {@link IODescriptor} - so it lives with the enum now.
|
||||
*/
|
||||
public static ExecutionVariableKind forDescriptor(IODescriptor descriptor) {
|
||||
return switch (descriptor.getType()) {
|
||||
case TEXT -> TEXT;
|
||||
case BOOLEAN -> BOOLEAN;
|
||||
case FILE, CSV -> FILE_PATH;
|
||||
case JSON -> JSON;
|
||||
case ANY -> ANY;
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ import org.springframework.data.domain.Page;
|
|||
import org.springframework.data.domain.Pageable;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.util.StringUtils;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
import org.springframework.web.server.ResponseStatusException;
|
||||
import org.springframework.http.HttpStatus;
|
||||
|
|
@ -645,15 +646,36 @@ public class ExecutionsService {
|
|||
* steps. Falls back to a normal start otherwise.
|
||||
*/
|
||||
public ExecutionObject startContainerChild(String childId, ContainerExecutionContext executionContext) {
|
||||
if (executionContext.simulationEnabled() && executionContext.simulationDescriptor() != null) {
|
||||
ExecutionObject child = getExecution(childId);
|
||||
if (child.isSimulationAvailable()) {
|
||||
return startSimulationExecution(childId, executionContext.simulationDescriptor());
|
||||
}
|
||||
// Every container path - generic, iterator, loop and loop guard - starts its child here, so
|
||||
// this is the one place that can name the container in the failure. Without it the parent
|
||||
// records "Execution with id <uuid> cannot start", about an execution the user never sees.
|
||||
ExecutionObject child = getExecution(childId);
|
||||
String reason = child.notStartableReason();
|
||||
if (reason != null) {
|
||||
throw new IllegalStateException(
|
||||
"The subflow of " + containerLabel(executionContext) + " cannot start: " + reason);
|
||||
}
|
||||
if (executionContext.simulationEnabled() && executionContext.simulationDescriptor() != null
|
||||
&& child.isSimulationAvailable()) {
|
||||
return startSimulationExecution(childId, executionContext.simulationDescriptor());
|
||||
}
|
||||
return startExecution(childId);
|
||||
}
|
||||
|
||||
/** The container's own name where it can be resolved, so the message points at the diagram. */
|
||||
private String containerLabel(ContainerExecutionContext executionContext) {
|
||||
try {
|
||||
var step = getExecution(executionContext.parentExecutionId()).getContext().getSteps()
|
||||
.get(executionContext.parentStepId());
|
||||
if (step != null && step.getNode() != null && StringUtils.hasText(step.getNode().getName())) {
|
||||
return "container '" + step.getNode().getName() + "'";
|
||||
}
|
||||
} catch (RuntimeException ignored) {
|
||||
// Naming the container is a courtesy; never let it replace the real failure.
|
||||
}
|
||||
return "this container";
|
||||
}
|
||||
|
||||
/**
|
||||
* Durable startup reconciliation (Q3/Q6): re-establishes the container
|
||||
* coordinator for every linked child, without relying on any in-memory
|
||||
|
|
@ -1060,23 +1082,8 @@ public class ExecutionsService {
|
|||
}
|
||||
}
|
||||
setExecutionVariableDescriptors(child.getId(), parent.getContext().getExecutionVariableDescriptors());
|
||||
// The parent's own global inputs, read from the map that holds them. This used to take a
|
||||
// "global." view of getExecutionVariables(), which is the *unprefixed* variables map and so
|
||||
// never matched the prefix: every child was handed null values, sat in CREATED and failed
|
||||
// the parent with "is not in READY status". The prefixed keys live in the runtime map,
|
||||
// which is deliberately not exposed - so go to the source instead.
|
||||
Map<String, Object> globalInputs = new LinkedHashMap<>(parent.getContext().getGlobalInputs());
|
||||
Map<String, ExecutionVariableDescriptor> globalDescriptors = new LinkedHashMap<>();
|
||||
for (IODescriptor required : child.getRequiredGlobalInputs()) {
|
||||
globalDescriptors.put(required.getName(), ExecutionVariableDescriptor.builder()
|
||||
.name(required.getName())
|
||||
.kind(ContainerExecutionSupport.mapExecutionVariableKind(required))
|
||||
.value(globalInputs.get(required.getName()))
|
||||
.cleanupPolicy(ExecutionVariableCleanupPolicy.NONE)
|
||||
.description("Global flow input")
|
||||
.build());
|
||||
}
|
||||
setGlobalInputDescriptors(child.getId(), globalDescriptors);
|
||||
setGlobalInputDescriptors(child.getId(),
|
||||
SubflowGlobalInputs.descriptorsFor(child, parent.getContext().getGlobalInputs()));
|
||||
for (var entry : inputPortsByName.entrySet()) {
|
||||
if (!inputValues.containsKey(entry.getKey())) {
|
||||
throw new IllegalArgumentException(containerTypeName + " subflow input is missing: " + entry.getKey());
|
||||
|
|
|
|||
|
|
@ -0,0 +1,44 @@
|
|||
package it.cnr.isti.workflow.manager.executions;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import it.cnr.isti.workflow.manager.ios.IODescriptor;
|
||||
|
||||
/**
|
||||
* Hands a container's subflow the global inputs its parent was given.
|
||||
*
|
||||
* <p>A subflow is its own execution, so the values the user typed on the parent do not reach it by
|
||||
* themselves. Every container path used to carry them across with its own copy of this code, and
|
||||
* the copies drifted: the one in {@code ExecutionsService} read the unprefixed variables map
|
||||
* through a {@code global.} view, matched nothing, and handed every iterator and loop child null
|
||||
* values - while {@code GenericContainerExecutor}, reading the runtime map, kept working. One copy,
|
||||
* so the next divergence cannot happen.
|
||||
*/
|
||||
public final class SubflowGlobalInputs {
|
||||
|
||||
private SubflowGlobalInputs() {
|
||||
}
|
||||
|
||||
/**
|
||||
* Descriptors for exactly the globals the child declares, valued from the parent's own global
|
||||
* inputs. A global the parent does not have stays null, which is what leaves the child unable
|
||||
* to start - reported by {@link ExecutionObject#notStartableReason()} rather than silently.
|
||||
*/
|
||||
public static Map<String, ExecutionVariableDescriptor> descriptorsFor(ExecutionObject child,
|
||||
Map<String, Object> parentGlobalInputs) {
|
||||
Map<String, Object> values = parentGlobalInputs == null ? Map.of() : parentGlobalInputs;
|
||||
Map<String, ExecutionVariableDescriptor> descriptors = new LinkedHashMap<>();
|
||||
for (IODescriptor required : child.getRequiredGlobalInputs()) {
|
||||
descriptors.put(required.getName(), ExecutionVariableDescriptor.builder()
|
||||
.name(required.getName())
|
||||
.kind(ExecutionVariableKind.forDescriptor(required))
|
||||
.multiple(required.isMultiple())
|
||||
.value(values.get(required.getName()))
|
||||
.cleanupPolicy(ExecutionVariableCleanupPolicy.NONE)
|
||||
.description("Global flow input")
|
||||
.build());
|
||||
}
|
||||
return descriptors;
|
||||
}
|
||||
}
|
||||
|
|
@ -24,6 +24,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionsService;
|
|||
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.bias.BiasExecutionContext;
|
||||
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasContainerPropagation;
|
||||
import it.cnr.isti.workflow.manager.executions.persistence.ContainerContinuationPhase;
|
||||
|
|
@ -70,17 +71,8 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
|
||||
propagateAuthorizations(innerExecution, authorizations);
|
||||
executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors);
|
||||
executionsService.setGlobalInputDescriptors(innerExecution.getId(),
|
||||
innerExecution.getRequiredGlobalInputs().stream()
|
||||
.collect(java.util.stream.Collectors.toMap(
|
||||
required -> required.getName(),
|
||||
required -> ExecutionVariableDescriptor.builder()
|
||||
.name(required.getName())
|
||||
.kind(mapKind(required))
|
||||
.value(ExecutionRuntimeContextSupport.globalView(executionVariables).get(required.getName()))
|
||||
.cleanupPolicy(it.cnr.isti.workflow.manager.executions.ExecutionVariableCleanupPolicy.NONE)
|
||||
.description("Global flow input")
|
||||
.build())));
|
||||
executionsService.setGlobalInputDescriptors(innerExecution.getId(), SubflowGlobalInputs.descriptorsFor(
|
||||
innerExecution, ExecutionRuntimeContextSupport.globalView(executionVariables)));
|
||||
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getSubFlow()).stream()
|
||||
|
|
@ -203,14 +195,4 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
return GenericContainerType.class;
|
||||
}
|
||||
|
||||
private it.cnr.isti.workflow.manager.executions.ExecutionVariableKind mapKind(
|
||||
it.cnr.isti.workflow.manager.ios.IODescriptor descriptor) {
|
||||
return switch (descriptor.getType()) {
|
||||
case TEXT -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.TEXT;
|
||||
case BOOLEAN -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.BOOLEAN;
|
||||
case FILE, CSV -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.FILE_PATH;
|
||||
case JSON -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.JSON;
|
||||
case ANY -> it.cnr.isti.workflow.manager.executions.ExecutionVariableKind.ANY;
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -148,7 +148,7 @@ class BranchRejoinConcurrencyTest {
|
|||
.<Runnable>mapToObj(ignored -> () -> executionsService.cancelExecution(execution.getId()))
|
||||
.toList();
|
||||
assertDoesNotThrow(() -> {
|
||||
List<java.util.concurrent.Future<?>> futures = cancels.stream().map(pool::submit).toList();
|
||||
List<? extends java.util.concurrent.Future<?>> futures = cancels.stream().map(pool::submit).toList();
|
||||
for (var future : futures) {
|
||||
future.get(10, TimeUnit.SECONDS);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1433,6 +1433,51 @@ public class ExecutionTest {
|
|||
finished.getContext().getResult().values().stream().findFirst().orElseThrow());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void aSubflowThatCannotStartNamesTheContainerAndWhatIsMissing() {
|
||||
// The subflow declares a global its parent does not, so nothing can supply it. The failure
|
||||
// has to say which container and which input: it used to read "Execution with id <uuid> is
|
||||
// not in READY status (CURRENT STATUS is CREATED)", naming an execution the user never sees.
|
||||
Block<LLMBlockType> internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder()
|
||||
.name("Internal LLM")
|
||||
.llmDescriptor(llmBrick)
|
||||
.prompt("Hello, ${{name}} from ${{global.who}}!")
|
||||
.build());
|
||||
|
||||
Container<IteratorContainerType> container = iteratorContainerFactory.create(IteratorContainerConfiguration.builder()
|
||||
.name("Candidate loop")
|
||||
.subFlow(FlowData.builder()
|
||||
.block(internalBlock)
|
||||
.globalInput(IODescriptor.input("who", IOType.TEXT, false, null))
|
||||
.build())
|
||||
.iterationInput("name")
|
||||
.build());
|
||||
|
||||
ExecutionObject execObject = executionsService.createExecution("Iterator flow",
|
||||
FlowData.builder().container(container).build());
|
||||
|
||||
executionsService.prepareInput(execObject.getId(), container.getId(), "name", List.of("Alice"));
|
||||
execObject = executionsService.startExecution(execObject.getId());
|
||||
for (int attempt = 0; attempt < 200 && !execObject.getContext().getStatus().isFinalState(); attempt++) {
|
||||
try {
|
||||
Thread.sleep(50);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
execObject = executionsService.getExecution(execObject.getId());
|
||||
}
|
||||
|
||||
ExecutionObject failed = execObject;
|
||||
assertEquals(ExecutionStatus.ERROR, failed.getContext().getStatus());
|
||||
String message = failed.getContext().getErrors().values().stream()
|
||||
.map(String::valueOf).findFirst().orElseThrow();
|
||||
assertTrue(message.contains("Candidate loop"), () -> "should name the container: " + message);
|
||||
assertTrue(message.contains("who"), () -> "should name the missing global input: " + message);
|
||||
// And the error is filed against the container step, so the diagram points at it too.
|
||||
assertTrue(failed.getContext().getErrors().containsKey(container.getId()));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void genericContainerSuspendsForInnerInteractionAndResumesFromChild() {
|
||||
Block<HumanInteractionBlockType> innerInteraction = humanInteractiveBlockFactory.create(
|
||||
|
|
|
|||
Loading…
Reference in New Issue