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(