From 2ea8e339873217cd70734023350c3a42f210dc12 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Thu, 3 Sep 2026 15:34:21 +0200 Subject: [PATCH] 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) --- .../executions/ContainerExecutionSupport.java | 11 ---- .../manager/executions/ExecutionObject.java | 60 ++++++++++++++----- .../executions/ExecutionVariableKind.java | 18 +++++- .../manager/executions/ExecutionsService.java | 51 +++++++++------- .../executions/SubflowGlobalInputs.java | 44 ++++++++++++++ .../containers/GenericContainerExecutor.java | 24 +------- .../BranchRejoinConcurrencyTest.java | 2 +- .../manager/executions/ExecutionTest.java | 45 ++++++++++++++ 8 files changed, 185 insertions(+), 70 deletions(-) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/SubflowGlobalInputs.java diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java index 5501d2e..d460774 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionSupport.java @@ -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 iteratorState(List remainingValues, Map runtimeInputValues) { Map state = new LinkedHashMap<>(); state.put("remainingValues", new java.util.ArrayList<>(remainingValues)); 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 228d285..3c4c48c 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 @@ -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. + * + *

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 reasons = new ArrayList<>(); + List missingGlobalInputs = getMissingGlobalInputKeys(); + if (!missingGlobalInputs.isEmpty()) { + reasons.add("no value for the global input" + (missingGlobalInputs.size() == 1 ? " " : "s ") + + quoted(missingGlobalInputs)); + } + List missingAuthorizations = getMissingAuthorizationKeys(); + if (!missingAuthorizations.isEmpty()) { + reasons.add("no credential for " + quoted(missingAuthorizations)); + } + List unsetInputs = new ArrayList<>(); + for (Step step : this.context.getSteps().values()) { + List 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 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; - }; - } - } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java index 693c324..7ab3388 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionVariableKind.java @@ -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; + }; + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java index b753365..2f851b2 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -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 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 globalInputs = new LinkedHashMap<>(parent.getContext().getGlobalInputs()); - Map 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()); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/SubflowGlobalInputs.java b/src/main/java/it/cnr/isti/workflow/manager/executions/SubflowGlobalInputs.java new file mode 100644 index 0000000..b27fec3 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/SubflowGlobalInputs.java @@ -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. + * + *

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 descriptorsFor(ExecutionObject child, + Map parentGlobalInputs) { + Map values = parentGlobalInputs == null ? Map.of() : parentGlobalInputs; + Map 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java index 8d10c58..fa0f7de 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/GenericContainerExecutor.java @@ -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 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 inputPortsByName = ContainerFlowInterfaceResolver .getExposedInputs(configuration.getSubFlow()).stream() @@ -203,14 +195,4 @@ public class GenericContainerExecutor implements ContainerExecutor 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; - }; - } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/BranchRejoinConcurrencyTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/BranchRejoinConcurrencyTest.java index 744c1a2..399f535 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/BranchRejoinConcurrencyTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/BranchRejoinConcurrencyTest.java @@ -148,7 +148,7 @@ class BranchRejoinConcurrencyTest { .mapToObj(ignored -> () -> executionsService.cancelExecution(execution.getId())) .toList(); assertDoesNotThrow(() -> { - List> futures = cancels.stream().map(pool::submit).toList(); + List> futures = cancels.stream().map(pool::submit).toList(); for (var future : futures) { future.get(10, TimeUnit.SECONDS); } 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 072e626..53c9ac6 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 @@ -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 is + // not in READY status (CURRENT STATUS is CREATED)", naming an execution the user never sees. + Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() + .name("Internal LLM") + .llmDescriptor(llmBrick) + .prompt("Hello, ${{name}} from ${{global.who}}!") + .build()); + + Container 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 innerInteraction = humanInteractiveBlockFactory.create(