diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionContext.java index 8773339..a1db64b 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ContainerExecutionContext.java @@ -1,7 +1,14 @@ package it.cnr.isti.workflow.manager.executions; -/** Stable identity of a parent execution and its active container step. */ -public record ContainerExecutionContext(String parentExecutionId, String parentStepId) { +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; + +/** + * Stable identity of a parent execution and its active container step, plus + * whether the parent execution is running in simulation mode so a container + * executor can propagate it to the child it creates. + */ +public record ContainerExecutionContext(String parentExecutionId, String parentStepId, boolean simulationEnabled, + LLMDescriptor simulationDescriptor) { public ContainerExecutionContext { if (parentExecutionId == null || parentExecutionId.isBlank()) { @@ -11,4 +18,8 @@ public record ContainerExecutionContext(String parentExecutionId, String parentS throw new IllegalArgumentException("Parent step id is required"); } } + + public ContainerExecutionContext(String parentExecutionId, String parentStepId) { + this(parentExecutionId, parentStepId, false, null); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java index cef3f7e..8d2be3a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java @@ -87,6 +87,7 @@ public class ExecutionContext implements ExecutionListener { protected void setInteractionSimulationEnabled(boolean interactionSimulationEnabled) { this.interactionSimulationEnabled = interactionSimulationEnabled; + this.steps.values().forEach(step -> step.setExecutionSimulationEnabled(interactionSimulationEnabled)); } protected void setInteractionSimulationDescriptor(LLMDescriptor interactionSimulationDescriptor) { @@ -677,7 +678,7 @@ public class ExecutionContext implements ExecutionListener { } this.startTime = snapshot.getStartTime(); this.endTime = snapshot.getEndTime(); - this.interactionSimulationEnabled = snapshot.isInteractionSimulationEnabled(); + this.setInteractionSimulationEnabled(snapshot.isInteractionSimulationEnabled()); this.setInteractionSimulationDescriptor(snapshot.getInteractionSimulationDescriptor()); this.setBiasExecutionContext(snapshot.getBiasExecutionContext()); refreshRuntimeExecutionVariables(); 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 3b8c9cb..47ceece 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 @@ -14,6 +14,8 @@ import java.util.stream.Collectors; import com.fasterxml.jackson.annotation.JsonIgnore; import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration; import it.cnr.isti.workflow.manager.executions.persistence.ExecutionSnapshot; import it.cnr.isti.workflow.manager.executions.persistence.ExecutionStepSnapshot; import it.cnr.isti.workflow.manager.executions.persistence.ContainerContinuationSnapshot; @@ -436,6 +438,17 @@ public class ExecutionObject { private boolean hasSimulationAvailable(List> steps) { boolean hasInteractiveSteps = false; for (Step step : steps) { + if (step.getNode() instanceof Container container) { + Boolean containerAvailability = containerSimulationAvailability(container); + if (containerAvailability == null) { + continue; + } + hasInteractiveSteps = true; + if (!containerAvailability) { + return false; + } + continue; + } if (!step.getNode().isUserInteractive()) { continue; } @@ -447,6 +460,39 @@ public class ExecutionObject { return hasInteractiveSteps; } + /** + * Whether a container's inner subflow(s) make simulation available, + * propagated from the container's own step so a top-level execution + * recognizes an interactive node nested inside a container. Containers + * are never themselves {@code isUserInteractive()}. Returns {@code null} + * when the subflow(s) have no interactive node at all (doesn't affect + * the parent's overall availability, matching a non-interactive step). + */ + private Boolean containerSimulationAvailability(Container container) { + List subFlows = new ArrayList<>(); + if (container.getSpecificConfiguration() instanceof ContainerConfiguration configuration + && configuration.getSubFlow() != null) { + subFlows.add(configuration.getSubFlow()); + } + if (container.getSpecificConfiguration() instanceof LoopContainerConfiguration loopConfiguration + && loopConfiguration.getGuardSubFlow() != null) { + subFlows.add(loopConfiguration.getGuardSubFlow()); + } + boolean hasInteractive = false; + for (FlowData subFlow : subFlows) { + for (FlowNode node : subFlow.getNodes()) { + if (!node.isUserInteractive()) { + continue; + } + hasInteractive = true; + if (!NodeExecutors.supportsSimulation(node)) { + return false; + } + } + } + return hasInteractive ? true : null; + } + protected void abortOnError() { if (this.executorService != null) { this.executorService.shutdownNow(); 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 b5fc651..8c10362 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 @@ -15,6 +15,8 @@ import jakarta.annotation.PostConstruct; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; import org.springframework.data.domain.Page; import org.springframework.data.domain.Pageable; import org.springframework.scheduling.annotation.Scheduled; @@ -38,10 +40,19 @@ import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfigurati import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver; +import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; +import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; import it.cnr.isti.workflow.manager.executions.persistence.ExecutionSnapshot; +import it.cnr.isti.workflow.manager.executions.persistence.ContainerContinuationPhase; import it.cnr.isti.workflow.manager.executions.persistence.ContainerContinuationSnapshot; import it.cnr.isti.workflow.manager.executions.steps.Step; +import it.cnr.isti.workflow.manager.executions.executors.BooleanLlmResponseParser; +import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.executions.bias.BiasActivation; import it.cnr.isti.workflow.manager.executions.bias.BiasApiException; import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext; @@ -485,25 +496,67 @@ public class ExecutionsService { } if (execution.getContext().getStatus() == ExecutionStatus.RUNNING) throw new IllegalStateException("Execution with id " + id + " is still running"); + // Children are deleted transitively by the DB FK (ON DELETE CASCADE); evict + // them from the in-memory cache too so a stale cached child isn't served. + List childIds = executionRepository.findByParentExecutionIdOrderByCreationTimeAsc(id).stream() + .map(ExecutionEntity::getId) + .toList(); execution.shutdown(); executions.remove(id); lastAccessByExecutionId.remove(id); biasImpactJobRepository.deleteByExecutionId(id); biasImpactReportRepository.deleteByBaselineExecutionIdOrBiasedExecutionId(id, id); executionRepository.deleteById(id); + for (String childId : childIds) { + ExecutionObject cachedChild = executions.remove(childId); + if (cachedChild != null) { + cachedChild.shutdown(); + } + lastAccessByExecutionId.remove(childId); + } } + /** + * Cancels an execution. If it is a container's active child, the parent + * container step is failed accordingly (the same outcome as an + * unexpected child error). If it has an active child of its own + * (a step in {@code WAITING_FOR_SUBFLOW}), that child is cancelled too. + * Both directions are idempotent: cancelling an already-final execution + * or an already-reconciled child is a no-op. + */ @Transactional public ExecutionObject cancelExecution(String id) { ExecutionObject eo = getExecution(id); if (eo.getContext().getStatus().isFinalState()) { return eo; } + List activeChildIds = activeContainerChildIds(eo); eo.cancel(); touchExecution(id); + for (String childId : activeChildIds) { + try { + cancelExecution(childId); + } catch (RuntimeException exception) { + logger.warn("Failed to cancel subflow child {} of execution {}", childId, id, exception); + } + } + if (eo.getExecutionKind() == ExecutionKind.SUBFLOW && eo.getParentExecutionId() != null + && eo.getParentStepId() != null) { + reconcileContainerSubflow(eo.getParentExecutionId(), eo.getParentStepId(), id); + } return eo; } + private List activeContainerChildIds(ExecutionObject execution) { + return execution.getContext().getSteps().values().stream() + .filter(step -> step.getStatus() == it.cnr.isti.workflow.manager.executions.steps.StepStatus.WAITING_FOR_SUBFLOW) + .map(Step::getContainerContinuation) + .filter(Objects::nonNull) + .map(ContainerContinuationSnapshot::getActiveInnerExecutionId) + .filter(Objects::nonNull) + .toList(); + } + public ExecutionObject prepareInput(String executionId, String blockId, String inputName, Object input) { ExecutionObject eo = getExecution(executionId); eo.setInput(blockId, inputName, input); @@ -524,18 +577,59 @@ public class ExecutionsService { } /** - * Arms the in-memory continuation used by the GenericContainer Tappa A - * runtime. The persisted continuation remains the source of identity; - * durable reconciliation after a restart is implemented in Tappa B. + * Starts a container's child execution, propagating simulation mode from + * the parent step when requested and the child actually has simulable + * steps. Falls back to a normal start otherwise. */ - public void watchGenericSubflow(String parentExecutionId, String parentStepId, String childExecutionId) { - ExecutionObject child = getExecution(childExecutionId); - child.getContext().addEventListener(ignored -> reconcileGenericSubflow( - parentExecutionId, parentStepId, childExecutionId)); - reconcileGenericSubflow(parentExecutionId, parentStepId, childExecutionId); + 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()); + } + } + return startExecution(childId); } - private synchronized void reconcileGenericSubflow(String parentExecutionId, String parentStepId, + /** + * Durable startup reconciliation (Q3/Q6): re-establishes the container + * coordinator for every linked child, without relying on any in-memory + * listener having survived a restart. A child that already reached a + * final state while the process was down is reconciled immediately; a + * still-pending child just gets its listener re-armed for whenever it + * eventually finishes. Safe to run repeatedly (reconciliation is + * idempotent) and failures on one child don't block the others. + */ + @Transactional + @EventListener(ApplicationReadyEvent.class) + public void reconcileContainerSubflowsOnStartup() { + for (ExecutionEntity child : executionRepository.findByExecutionKind(ExecutionKind.SUBFLOW)) { + if (child.getParentExecutionId() == null || child.getParentStepId() == null) { + continue; + } + try { + watchContainerSubflow(child.getParentExecutionId(), child.getParentStepId(), child.getId()); + } catch (RuntimeException exception) { + logger.warn("Failed to reconcile subflow child {} during startup", child.getId(), exception); + } + } + } + + /** + * Arms the durable coordinator for a container's active child: attaches + * an in-memory listener for immediate notifications (an optimization) + * and immediately reconciles once, so a child that already reached a + * final state (e.g. discovered during startup reconciliation) is picked + * up right away regardless of whether the listener ever fires. + */ + public void watchContainerSubflow(String parentExecutionId, String parentStepId, String childExecutionId) { + ExecutionObject child = getExecution(childExecutionId); + child.getContext().addEventListener(ignored -> reconcileContainerSubflow( + parentExecutionId, parentStepId, childExecutionId)); + reconcileContainerSubflow(parentExecutionId, parentStepId, childExecutionId); + } + + private synchronized void reconcileContainerSubflow(String parentExecutionId, String parentStepId, String childExecutionId) { ExecutionObject child = getExecution(childExecutionId); if (!child.getContext().getStatus().isFinalState()) { @@ -543,32 +637,514 @@ public class ExecutionsService { } ExecutionObject parent = getExecution(parentExecutionId); Step parentStep = parent.getContext().getSteps().get(parentStepId); - if (parentStep == null || !(parentStep.getNode() instanceof Container container) - || !(container.getSpecificConfiguration() instanceof GenericContainerConfiguration configuration)) { + if (parentStep == null || !(parentStep.getNode() instanceof Container container)) { return; } ContainerContinuationSnapshot continuation = parentStep.getContainerContinuation(); if (continuation == null || !childExecutionId.equals(continuation.getActiveInnerExecutionId())) { return; } + Object configuration = container.getSpecificConfiguration(); + if (configuration instanceof GenericContainerConfiguration generic) { + reconcileGenericSubflow(parent, parentStep, generic, child); + } else if (configuration instanceof LoopContainerConfiguration loop) { + reconcileLoopSubflow(parent, parentStep, container, loop, child, continuation); + } else if (configuration instanceof IteratorContainerConfiguration iterator) { + reconcileIteratorSubflow(parent, parentStep, container, iterator, child, continuation); + } + } + + private void reconcileGenericSubflow(ExecutionObject parent, Step parentStep, + GenericContainerConfiguration configuration, ExecutionObject child) { if (child.getContext().getStatus() == ExecutionStatus.SUCCESS) { parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors()); - parent.completeContainerSubflow(parentStepId, - NodeExecutionResult.completed(resolveGenericContainerOutputs(child, configuration))); + parent.completeContainerSubflow(parentStep.getId(), + NodeExecutionResult.completed(collectExposedOutputsAsMap(child, + ContainerFlowInterfaceResolver.getExposedOutputs(configuration.getSubFlow())))); return; } - parent.failContainerSubflow(parentStepId, + parent.failContainerSubflow(parentStep.getId(), "GenericContainer subflow ended with status " + child.getContext().getStatus() + (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors())); } - private Map resolveGenericContainerOutputs(ExecutionObject child, - GenericContainerConfiguration configuration) { + // ---- Iterator: sequential per-item children, no guard ---- + + /** + * Entry point used by {@code IteratorContainerExecutor.execute(...)} to + * drive the (possibly empty) list of iterations. Iterations that finish + * synchronously are chained without ever suspending the step; the + * moment one doesn't, the step suspends and the container is resumed + * later by {@link #reconcileIteratorSubflow}. + */ + public NodeExecutionResult startIteratorSubflow(String parentExecutionId, String parentStepId, + Container container, IteratorContainerConfiguration configuration, + IteratorContainerInterfaceResolver.Resolution resolution, List iterationValues, + Map runtimeInputValues, Map authorizations, ExecutionEventLogger eventLogger, + BiasExecutionContext innerBiasExecutionContext, ContainerExecutionContext executionContext) { + Map> accumulated = new LinkedHashMap<>(); + for (var output : resolution.resolvedOutputs()) { + accumulated.put(output.publicName(), new java.util.ArrayList<>()); + } + if (iterationValues.isEmpty()) { + return NodeExecutionResult.completed(widen(accumulated)); + } + ExecutionObject parent = getExecution(parentExecutionId); + List remaining = new java.util.ArrayList<>(iterationValues.subList(1, iterationValues.size())); + ContainerAdvanceOutcome outcome = runIteratorIterations(parent, parentStepId, container, configuration, resolution, + 1, iterationValues.getFirst(), remaining, runtimeInputValues, accumulated, authorizations, eventLogger, + innerBiasExecutionContext, executionContext); + return toNodeExecutionResult(outcome, parentExecutionId, parentStepId); + } + + private void reconcileIteratorSubflow(ExecutionObject parent, Step parentStep, Container container, + IteratorContainerConfiguration configuration, ExecutionObject child, ContainerContinuationSnapshot continuation) { + if (child.getContext().getStatus() != ExecutionStatus.SUCCESS) { + parent.failContainerSubflow(parentStep.getId(), + "IteratorContainer subflow ended with status " + child.getContext().getStatus() + + (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors())); + return; + } + parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors()); + IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration); + List outputHandles = resolution.resolvedOutputs().stream() + .map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList(); + Map iterationOutputs = collectExposedOutputsAsMap(child, outputHandles); + Map> accumulated = asAccumulatedOutputs(continuation.getAccumulatedOutputs()); + for (var output : resolution.resolvedOutputs()) { + accumulated.computeIfAbsent(output.publicName(), ignored -> new java.util.ArrayList<>()) + .add(iterationOutputs.get(output.publicName())); + } + int completedIndex = continuation.getIterationIndex() == null ? 1 : continuation.getIterationIndex(); + ExecutionEventLogger eventLogger = parentStep.getEventLogger(); + if (eventLogger != null) { + eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, "Completed iterator iteration " + completedIndex, + Map.of("containerType", IteratorContainerType.TYPE, "iteration", completedIndex)); + } + Map state = continuation.getState() == null ? Map.of() : continuation.getState(); + List remaining = asObjectList(state.get("remainingValues")); + Map runtimeInputValues = asStringObjectMap(state.get("runtimeInputValues")); + + if (remaining.isEmpty()) { + setContainerContinuation(parent.getId(), parentStep.getId(), null); + parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(widen(accumulated))); + return; + } + Object nextValue = remaining.getFirst(); + List nextRemaining = new java.util.ArrayList<>(remaining.subList(1, remaining.size())); + BiasExecutionContext innerBiasExecutionContext = BiasContainerPropagation.innerContextFor(container.getId(), + List.of(configuration.getSubFlow()), parent.getBiasExecutionContext()); + ContainerExecutionContext executionContext = new ContainerExecutionContext(parent.getId(), parentStep.getId(), + parentStep.isExecutionSimulationEnabled(), parentStep.getInteractionSimulationDescriptor()); + ContainerAdvanceOutcome outcome = runIteratorIterations(parent, parentStep.getId(), container, configuration, + resolution, completedIndex + 1, nextValue, nextRemaining, runtimeInputValues, accumulated, + parent.getProvidedAuthorizations(), eventLogger, innerBiasExecutionContext, executionContext); + if (outcome.completed()) { + parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(outcome.outputs())); + } else { + watchContainerSubflow(parent.getId(), parentStep.getId(), outcome.pendingChildId()); + } + } + + private ContainerAdvanceOutcome runIteratorIterations(ExecutionObject parent, String parentStepId, + Container container, IteratorContainerConfiguration configuration, + IteratorContainerInterfaceResolver.Resolution resolution, int startIterationIndex, Object startValue, + List remainingValues, Map runtimeInputValues, Map> accumulatedOutputs, + Map authorizations, ExecutionEventLogger eventLogger, BiasExecutionContext innerBiasExecutionContext, + ContainerExecutionContext executionContext) { + int iterationIndex = startIterationIndex; + Object iterationValue = startValue; + List remaining = remainingValues; + Map inputPortsByName = resolution.resolvedInputs().stream() + .collect(Collectors.toMap(IteratorContainerInterfaceResolver.ResolvedInput::publicName, + IteratorContainerInterfaceResolver.ResolvedInput::exposedHandle)); + List outputHandles = resolution.resolvedOutputs().stream() + .map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList(); + + while (true) { + if (eventLogger != null) { + eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED, "Starting iterator iteration " + iterationIndex, + Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex)); + } + Map inputValues = new LinkedHashMap<>(runtimeInputValues); + inputValues.put(resolution.iteratedPublicName(), iterationValue); + + ExecutionObject child = createAndStartSubflowChild(parent, parentStepId, container, + container.getName() + " iteration " + iterationIndex, configuration.getSubFlow(), inputPortsByName, + inputValues, iterationIndex, ContainerSubflowRole.MAIN, authorizations, eventLogger, + IteratorContainerType.TYPE, innerBiasExecutionContext, executionContext, + iteratorState(remaining, runtimeInputValues), widen(accumulatedOutputs)); + + if (!child.getContext().getStatus().isFinalState()) { + return ContainerAdvanceOutcome.suspended(child.getId()); + } + if (child.getContext().getStatus() != ExecutionStatus.SUCCESS) { + throw new IllegalStateException("IteratorContainer subflow ended with status " + child.getContext().getStatus() + + (child.getContext().getErrors().isEmpty() ? "" : ": " + child.getContext().getErrors())); + } + parent.setExecutionVariableDescriptors(child.getContext().getExecutionVariableDescriptors()); + Map iterationOutputs = collectExposedOutputsAsMap(child, outputHandles); + for (var output : resolution.resolvedOutputs()) { + accumulatedOutputs.get(output.publicName()).add(iterationOutputs.get(output.publicName())); + } + if (eventLogger != null) { + eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, "Completed iterator iteration " + iterationIndex, + Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex)); + } + if (remaining.isEmpty()) { + setContainerContinuation(parent.getId(), parentStepId, null); + return ContainerAdvanceOutcome.completed(widen(accumulatedOutputs)); + } + iterationValue = remaining.getFirst(); + remaining = new java.util.ArrayList<>(remaining.subList(1, remaining.size())); + iterationIndex++; + } + } + + // ---- Loop: main subflow + guard subflow per iteration ---- + + private enum LoopPhase { MAIN, GUARD } + + /** Entry point used by {@code LoopContainerExecutor.execute(...)}. */ + public NodeExecutionResult startLoopSubflow(String parentExecutionId, String parentStepId, Container container, + LoopContainerConfiguration configuration, Map initialInputs, Map authorizations, + ExecutionEventLogger eventLogger, BiasExecutionContext innerBiasExecutionContext, + ContainerExecutionContext executionContext) { + ExecutionObject parent = getExecution(parentExecutionId); + ContainerAdvanceOutcome outcome = runLoopFrom(parent, parentStepId, container, configuration, LoopPhase.MAIN, 1, + initialInputs, Map.of(), authorizations, eventLogger, innerBiasExecutionContext, executionContext); + return toNodeExecutionResult(outcome, parentExecutionId, parentStepId); + } + + private void reconcileLoopSubflow(ExecutionObject parent, Step parentStep, Container container, + LoopContainerConfiguration configuration, ExecutionObject completedChild, ContainerContinuationSnapshot continuation) { + if (completedChild.getContext().getStatus() != ExecutionStatus.SUCCESS) { + parent.failContainerSubflow(parentStep.getId(), + "LoopContainer" + (completedChild.getSubflowRole() == ContainerSubflowRole.GUARD ? " guard" : "") + + " subflow ended with status " + completedChild.getContext().getStatus() + + (completedChild.getContext().getErrors().isEmpty() ? "" : ": " + completedChild.getContext().getErrors())); + return; + } + parent.setExecutionVariableDescriptors(completedChild.getContext().getExecutionVariableDescriptors()); + + Map state = continuation.getState() == null ? Map.of() : continuation.getState(); + Map currentInputs = asStringObjectMap(state.get("currentInputs")); + Map latestOutputs = asStringObjectMap(state.get("latestOutputs")); + int iterationIndex = continuation.getIterationIndex() == null ? 1 : continuation.getIterationIndex(); + ExecutionEventLogger eventLogger = parentStep.getEventLogger(); + BiasExecutionContext innerBiasExecutionContext = BiasContainerPropagation.innerContextFor(container.getId(), + List.of(configuration.getSubFlow(), configuration.getGuardSubFlow()), parent.getBiasExecutionContext()); + ContainerExecutionContext executionContext = new ContainerExecutionContext(parent.getId(), parentStep.getId(), + parentStep.isExecutionSimulationEnabled(), parentStep.getInteractionSimulationDescriptor()); + Map authorizations = parent.getProvidedAuthorizations(); + + ContainerAdvanceOutcome outcome; + if (completedChild.getSubflowRole() != ContainerSubflowRole.GUARD) { + List mainOutputHandles = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getSubFlow()); + Map freshLatest = collectExposedOutputsAsMap(completedChild, mainOutputHandles); + outcome = runLoopFrom(parent, parentStep.getId(), container, configuration, LoopPhase.GUARD, iterationIndex, + currentInputs, freshLatest, authorizations, eventLogger, innerBiasExecutionContext, executionContext); + } else { + Map inputPortsByName = ContainerFlowInterfaceResolver + .getExposedInputs(configuration.getSubFlow()).stream() + .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle)); + String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); + List guardOutputHandles = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getGuardSubFlow()); + Map guardResult = collectExposedOutputsAsMap(completedChild, guardOutputHandles); + boolean shouldContinue = parseLoopGuardResponse( + String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT))); + String feedback = String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT)); + logLoopGuardEvaluation(eventLogger, iterationIndex, shouldContinue); + if (!shouldContinue) { + setContainerContinuation(parent.getId(), parentStep.getId(), null); + parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(latestOutputs)); + return; + } + Map nextInputs = new LinkedHashMap<>(currentInputs); + nextInputs.put(feedbackInput, feedback); + outcome = runLoopFrom(parent, parentStep.getId(), container, configuration, LoopPhase.MAIN, iterationIndex + 1, + nextInputs, latestOutputs, authorizations, eventLogger, innerBiasExecutionContext, executionContext); + } + if (outcome.completed()) { + parent.completeContainerSubflow(parentStep.getId(), NodeExecutionResult.completed(outcome.outputs())); + } else { + watchContainerSubflow(parent.getId(), parentStep.getId(), outcome.pendingChildId()); + } + } + + private ContainerAdvanceOutcome runLoopFrom(ExecutionObject parent, String parentStepId, Container container, + LoopContainerConfiguration configuration, LoopPhase startPhase, int startIterationIndex, + Map startInputs, Map startLatestOutputs, Map authorizations, + ExecutionEventLogger eventLogger, BiasExecutionContext innerBiasExecutionContext, + ContainerExecutionContext executionContext) { + Map inputPortsByName = ContainerFlowInterfaceResolver + .getExposedInputs(configuration.getSubFlow()).stream() + .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle)); + String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); + List mainOutputHandles = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getSubFlow()); + Map guardInputsByName = ContainerFlowInterfaceResolver + .getExposedInputs(configuration.getGuardSubFlow()).stream() + .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle)); + List guardOutputHandles = ContainerFlowInterfaceResolver + .getExposedOutputs(configuration.getGuardSubFlow()); + + LoopPhase phase = startPhase; + int iterationIndex = startIterationIndex; + Map inputs = startInputs; + Map latestOutputs = startLatestOutputs == null ? Map.of() : startLatestOutputs; + + while (true) { + if (phase == LoopPhase.MAIN) { + if (iterationIndex > configuration.getMaxIterations()) { + throw new IllegalStateException( + "LoopContainer guard was not satisfied within maxIterations=" + configuration.getMaxIterations()); + } + if (eventLogger != null) { + eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED, "Starting loop iteration " + iterationIndex, + Map.of("containerType", LoopContainerType.TYPE, "iteration", iterationIndex)); + } + ExecutionObject mainChild = createAndStartSubflowChild(parent, parentStepId, container, + container.getName() + " iteration " + iterationIndex, configuration.getSubFlow(), inputPortsByName, + inputs, iterationIndex, ContainerSubflowRole.MAIN, authorizations, eventLogger, + LoopContainerType.TYPE, innerBiasExecutionContext, executionContext, + loopState(inputs, latestOutputs), null); + if (!mainChild.getContext().getStatus().isFinalState()) { + return ContainerAdvanceOutcome.suspended(mainChild.getId()); + } + if (mainChild.getContext().getStatus() != ExecutionStatus.SUCCESS) { + throw new IllegalStateException("LoopContainer subflow ended with status " + mainChild.getContext().getStatus() + + (mainChild.getContext().getErrors().isEmpty() ? "" : ": " + mainChild.getContext().getErrors())); + } + parent.setExecutionVariableDescriptors(mainChild.getContext().getExecutionVariableDescriptors()); + latestOutputs = collectExposedOutputsAsMap(mainChild, mainOutputHandles); + phase = LoopPhase.GUARD; + continue; + } + + Map guardTemplateValues = buildGuardTemplateValues(inputs, latestOutputs, iterationIndex); + Map nextInputsIfContinuing = new LinkedHashMap<>(inputs); + updateInputsForNextIteration(nextInputsIfContinuing, latestOutputs, inputPortsByName); + + ExecutionObject guardChild = createAndStartSubflowChild(parent, parentStepId, container, + container.getName() + " guard iteration " + iterationIndex, configuration.getGuardSubFlow(), + guardInputsByName, guardTemplateValues, iterationIndex, ContainerSubflowRole.GUARD, authorizations, + eventLogger, "LoopContainerGuard", innerBiasExecutionContext, executionContext, + loopState(nextInputsIfContinuing, latestOutputs), null); + if (!guardChild.getContext().getStatus().isFinalState()) { + return ContainerAdvanceOutcome.suspended(guardChild.getId()); + } + if (guardChild.getContext().getStatus() != ExecutionStatus.SUCCESS) { + throw new IllegalStateException("LoopContainer guard subflow ended with status " + guardChild.getContext().getStatus() + + (guardChild.getContext().getErrors().isEmpty() ? "" : ": " + guardChild.getContext().getErrors())); + } + parent.setExecutionVariableDescriptors(guardChild.getContext().getExecutionVariableDescriptors()); + Map guardResult = collectExposedOutputsAsMap(guardChild, guardOutputHandles); + boolean shouldContinue = parseLoopGuardResponse( + String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT))); + String feedback = String.valueOf(requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT)); + logLoopGuardEvaluation(eventLogger, iterationIndex, shouldContinue); + if (!shouldContinue) { + setContainerContinuation(parent.getId(), parentStepId, null); + return ContainerAdvanceOutcome.completed(new LinkedHashMap<>(latestOutputs)); + } + Map nextInputs = new LinkedHashMap<>(nextInputsIfContinuing); + nextInputs.put(feedbackInput, feedback); + inputs = nextInputs; + iterationIndex++; + phase = LoopPhase.MAIN; + } + } + + private void logLoopGuardEvaluation(ExecutionEventLogger eventLogger, int iterationIndex, boolean shouldContinue) { + if (eventLogger == null) { + return; + } + eventLogger.info(ExecutionEventType.CONTAINER_CONDITION_EVALUATED, "Evaluated loop guard", + Map.of("containerType", LoopContainerType.TYPE, "iteration", iterationIndex, "continue", shouldContinue)); + eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, + shouldContinue ? "Completed loop iteration " + iterationIndex : "Loop completed at iteration " + iterationIndex, + Map.of("containerType", LoopContainerType.TYPE, "iteration", iterationIndex, "continue", shouldContinue)); + } + + private static String resolveFeedbackInput(LoopContainerConfiguration configuration, + Map inputPortsByName) { + if (configuration.getFeedbackInput() != null && !configuration.getFeedbackInput().isBlank()) { + String configuredInput = configuration.getFeedbackInput(); + ContainerFlowInterfaceResolver.ExposedHandle handle = inputPortsByName.get(configuredInput); + if (handle == null || handle.handle().io().isMultiple()) { + throw new IllegalArgumentException( + "LoopContainer feedbackInput must target an open non-multiple subFlow input: " + configuredInput); + } + return configuredInput; + } + List eligibleInputs = inputPortsByName.values().stream() + .filter(handle -> !handle.handle().io().isMultiple()) + .toList(); + if (eligibleInputs.size() == 1) { + return eligibleInputs.getFirst().publicName(); + } + if (eligibleInputs.isEmpty()) { + throw new IllegalArgumentException( + "LoopContainer subFlow must expose at least one open non-multiple input to receive guard feedback"); + } + throw new IllegalArgumentException( + "LoopContainer feedbackInput is required when subFlow exposes more than one open non-multiple input"); + } + + private static void updateInputsForNextIteration(Map currentInputs, Map latestOutputs, + Map inputPortsByName) { + for (Map.Entry output : latestOutputs.entrySet()) { + if (inputPortsByName.containsKey(output.getKey())) { + currentInputs.put(output.getKey(), output.getValue()); + } + } + } + + private static Map buildGuardTemplateValues(Map currentInputs, + Map latestOutputs, int iteration) { + Map values = new LinkedHashMap<>(); + if (currentInputs != null) { + values.putAll(currentInputs); + currentInputs.forEach((key, value) -> values.put("inputs." + key, value)); + } + if (latestOutputs != null) { + values.putAll(latestOutputs); + latestOutputs.forEach((key, value) -> values.put("outputs." + key, value)); + } + values.put("iteration", iteration); + return values; + } + + private static Object requireGuardOutput(Map guardResult, String outputName) { + Object value = guardResult.get(outputName); + if (value == null) { + throw new IllegalArgumentException("LoopContainer guardSubFlow must return output: " + outputName); + } + return value; + } + + private static boolean parseLoopGuardResponse(String response) { + return BooleanLlmResponseParser.parse(response, "loop guard subflow", "Loop"); + } + + // ---- Shared container-child helpers ---- + + private NodeExecutionResult toNodeExecutionResult(ContainerAdvanceOutcome outcome, String parentExecutionId, + String parentStepId) { + if (outcome.completed()) { + return NodeExecutionResult.completed(outcome.outputs()); + } + String childId = outcome.pendingChildId(); + return NodeExecutionResult.suspended(() -> watchContainerSubflow(parentExecutionId, parentStepId, childId)); + } + + /** + * Creates a container child bound to the given subflow/inputs, persists + * the continuation before and after starting it (so the relationship is + * durable even if the process dies mid-start), and returns it either + * final or suspended. + */ + private ExecutionObject createAndStartSubflowChild(ExecutionObject parent, String parentStepId, Container container, + String childName, FlowData subFlow, Map inputPortsByName, + Map inputValues, int iterationIndex, ContainerSubflowRole role, Map authorizations, + ExecutionEventLogger eventLogger, String containerTypeName, BiasExecutionContext innerBiasExecutionContext, + ContainerExecutionContext executionContext, Map continuationState, + Map continuationAccumulatedOutputs) { + ExecutionObject child = createInnerExecution(childName, subFlow, innerBiasExecutionContext, parent, parentStepId, + iterationIndex, role); + forwardContainerChildEvents(child, eventLogger, iterationIndex, containerTypeName); + for (String key : child.getRequiredAuthorizations().stream().map(ExecutionAuthorizationRequirement::key).distinct().toList()) { + if (authorizations.containsKey(key)) { + setAuthorizationValue(child.getId(), key, authorizations.get(key)); + } + } + setExecutionVariableDescriptors(child.getId(), parent.getContext().getExecutionVariableDescriptors()); + Map globalInputs = ExecutionRuntimeContextSupport.globalView(parent.getContext().getExecutionVariables()); + Map globalDescriptors = new LinkedHashMap<>(); + for (IODescriptor required : child.getRequiredGlobalInputs()) { + globalDescriptors.put(required.getName(), ExecutionVariableDescriptor.builder() + .name(required.getName()) + .kind(mapExecutionVariableKind(required)) + .value(globalInputs.get(required.getName())) + .cleanupPolicy(ExecutionVariableCleanupPolicy.NONE) + .description("Global flow input") + .build()); + } + setGlobalInputDescriptors(child.getId(), globalDescriptors); + for (var entry : inputPortsByName.entrySet()) { + if (!inputValues.containsKey(entry.getKey())) { + throw new IllegalArgumentException(containerTypeName + " subflow input is missing: " + entry.getKey()); + } + var handle = entry.getValue().handle(); + prepareInput(child.getId(), handle.blockId(), handle.io().getName(), inputValues.get(entry.getKey())); + } + setContainerContinuation(parent.getId(), parentStepId, ContainerContinuationSnapshot.builder() + .activeInnerExecutionId(child.getId()) + .phase(ContainerContinuationPhase.CHILD_CREATED) + .iterationIndex(iterationIndex) + .state(continuationState) + .accumulatedOutputs(continuationAccumulatedOutputs == null ? Map.of() : continuationAccumulatedOutputs) + .revision(1) + .build()); + child = startContainerChild(child.getId(), executionContext); + setContainerContinuation(parent.getId(), parentStepId, ContainerContinuationSnapshot.builder() + .activeInnerExecutionId(child.getId()) + .phase(child.getContext().getStatus().isFinalState() ? ContainerContinuationPhase.CHILD_COMPLETED + : child.getContext().getStatus() == ExecutionStatus.WAITING ? ContainerContinuationPhase.WAITING_FOR_SUBFLOW + : ContainerContinuationPhase.CHILD_RUNNING) + .iterationIndex(iterationIndex) + .state(continuationState) + .accumulatedOutputs(continuationAccumulatedOutputs == null ? Map.of() : continuationAccumulatedOutputs) + .revision(2) + .build()); + return child; + } + + private void forwardContainerChildEvents(ExecutionObject child, ExecutionEventLogger eventLogger, Integer iterationIndex, + String containerTypeName) { + if (eventLogger == null) { + return; + } + child.getContext().addEventListener(event -> eventLogger.log( + event.getLevel(), + event.getType(), + event.getNodeName() == null || event.getNodeName().isBlank() ? event.getMessage() + : "[" + event.getNodeName() + "] " + event.getMessage(), + enrichChildEventDetails(event, child.getId(), iterationIndex, containerTypeName))); + } + + private static Map enrichChildEventDetails(ExecutionEvent event, String innerExecutionId, + Integer iterationIndex, String containerTypeName) { + Map details = new LinkedHashMap<>(); + details.put("containerType", containerTypeName); + details.put("innerExecutionId", innerExecutionId); + if (iterationIndex != null) { + details.put("iterationIndex", iterationIndex); + } + if (event.getNodeId() != null) { + details.put("innerNodeId", event.getNodeId()); + } + if (event.getNodeName() != null) { + details.put("innerNodeName", event.getNodeName()); + } + if (event.getStepId() != null) { + details.put("innerStepId", event.getStepId()); + } + if (event.getDetails() != null && !event.getDetails().isEmpty()) { + details.putAll(event.getDetails()); + } + return details; + } + + private static Map collectExposedOutputsAsMap(ExecutionObject execution, + List exposedOutputs) { Map outputs = new LinkedHashMap<>(); - for (var exposedHandle : it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver - .getExposedOutputs(configuration.getSubFlow())) { + for (var exposedHandle : exposedOutputs) { var handle = exposedHandle.handle(); - Object value = child.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); + Object value = execution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); if (value != null) { outputs.put(exposedHandle.publicName(), value); } @@ -576,6 +1152,68 @@ public class ExecutionsService { return outputs; } + private 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; + }; + } + + private static Map iteratorState(List remainingValues, Map runtimeInputValues) { + Map state = new LinkedHashMap<>(); + state.put("remainingValues", new java.util.ArrayList<>(remainingValues)); + state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues)); + return state; + } + + private static Map loopState(Map currentInputs, Map latestOutputs) { + Map state = new LinkedHashMap<>(); + state.put("currentInputs", new LinkedHashMap<>(currentInputs)); + state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs)); + return state; + } + + private static Map widen(Map> accumulated) { + return new LinkedHashMap<>(accumulated); + } + + @SuppressWarnings("unchecked") + private static List asObjectList(Object value) { + if (value instanceof List list) { + return new java.util.ArrayList<>((List) list); + } + return new java.util.ArrayList<>(); + } + + @SuppressWarnings("unchecked") + private static Map asStringObjectMap(Object value) { + if (value instanceof Map map) { + return new LinkedHashMap<>((Map) map); + } + return new LinkedHashMap<>(); + } + + private static Map> asAccumulatedOutputs(Map raw) { + Map> result = new LinkedHashMap<>(); + if (raw != null) { + raw.forEach((key, value) -> result.put(key, asObjectList(value))); + } + return result; + } + + private record ContainerAdvanceOutcome(boolean completed, Map outputs, String pendingChildId) { + static ContainerAdvanceOutcome completed(Map outputs) { + return new ContainerAdvanceOutcome(true, outputs, null); + } + + static ContainerAdvanceOutcome suspended(String pendingChildId) { + return new ContainerAdvanceOutcome(false, null, pendingChildId); + } + } + public ExecutionObject setAuthorizationValue(String executionId, String key, Object value) { ExecutionObject eo = getExecution(executionId); boolean knownRequirement = eo.getRequiredAuthorizations().stream() 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 6623f7f..1d13188 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 @@ -94,7 +94,7 @@ public class GenericContainerExecutor implements ContainerExecutor executionsService.watchGenericSubflow( + return NodeExecutionResult.suspended(() -> executionsService.watchContainerSubflow( executionContext.parentExecutionId(), executionContext.parentStepId(), childId)); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java index fbd5896..231f5d2 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/IteratorContainerExecutor.java @@ -1,53 +1,38 @@ package it.cnr.isti.workflow.manager.executions.executors.containers; -import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; -import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.Timer; - import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; import it.cnr.isti.workflow.manager.containers.iresolvers.IteratorContainerInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; -import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ContainerExecutionContext; -import it.cnr.isti.workflow.manager.executions.NodeExecutionResult; -import it.cnr.isti.workflow.manager.executions.ExecutionEvent; -import it.cnr.isti.workflow.manager.executions.ExecutionEventType; -import it.cnr.isti.workflow.manager.executions.ExecutionObject; -import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport; -import it.cnr.isti.workflow.manager.executions.ExecutionStatus; +import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; import it.cnr.isti.workflow.manager.executions.ExecutionsService; -import it.cnr.isti.workflow.manager.executions.FieldKey; +import it.cnr.isti.workflow.manager.executions.NodeExecutionResult; 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.steps.Input; +/** + * Iterates the subflow once per item of an input list. Iterations that + * complete synchronously (no interactive node reached) are chained without + * suspending the step; the moment one doesn't, the step suspends and is + * resumed later by {@link ExecutionsService}'s durable container coordinator. + */ @Component public class IteratorContainerExecutor implements ContainerExecutor { - private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(IteratorContainerExecutor.class); - private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; - private final ExecutionsService executionsService; - private final long innerExecutionWaitIntervalMs; - @Autowired(required = false) - MeterRegistry meterRegistry; - - public IteratorContainerExecutor( - ExecutionsService executionsService, - @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { + public IteratorContainerExecutor(ExecutionsService executionsService) { this.executionsService = executionsService; - this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -59,8 +44,7 @@ public class IteratorContainerExecutor implements ContainerExecutor input.getDescriptor().getName().equals(resolution.iteratedPublicName())) .findFirst() @@ -71,183 +55,24 @@ public class IteratorContainerExecutor implements ContainerExecutor runtimeInputsByName = inputs.stream() - .collect(LinkedHashMap::new, (map, input) -> map.put(input.getDescriptor().getName(), input), Map::putAll); - - Map> collectedOutputs = new LinkedHashMap<>(); - for (IteratorContainerInterfaceResolver.ResolvedOutput output : resolution.resolvedOutputs()) { - collectedOutputs.put(output.publicName(), new ArrayList<>()); - } - - int iterationIndex = 0; - for (Object iterationValue : iterationValues) { - iterationIndex++; - if (eventLogger != null) { - eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED, - "Starting iterator iteration " + iterationIndex, - Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex, "value", iterationValue)); - } - ExecutionObject innerExecution = executionsService.createInnerExecution(container.getName() + " iteration", - configuration.getSubFlow(), innerBiasExecutionContext); - forwardInnerEvents(innerExecution, eventLogger, iterationIndex, IteratorContainerType.TYPE); - - propagateAuthorizations(innerExecution, authorizations); - executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors); - Map globalInputs = ExecutionRuntimeContextSupport.globalView(executionVariables); - executionsService.setGlobalInputDescriptors(innerExecution.getId(), - innerExecution.getRequiredGlobalInputs().stream() - .collect(LinkedHashMap::new, (map, required) -> map.put(required.getName(), - ExecutionVariableDescriptor.builder() - .name(required.getName()) - .kind(mapKind(required)) - .value(globalInputs.get(required.getName())) - .cleanupPolicy(it.cnr.isti.workflow.manager.executions.ExecutionVariableCleanupPolicy.NONE) - .description("Global flow input") - .build()), Map::putAll)); - - for (IteratorContainerInterfaceResolver.ResolvedInput resolvedInput : resolution.resolvedInputs()) { - Object value = resolvedInput.iterated() - ? iterationValue - : runtimeInputsByName.get(resolvedInput.publicName()).getValue(); - executionsService.prepareInput( - innerExecution.getId(), - resolvedInput.exposedHandle().handle().blockId(), - resolvedInput.exposedHandle().handle().io().getName(), - value); - } - - innerExecution = startAndWait(innerExecution); - executionVariables.clear(); - executionVariables.putAll(innerExecution.getContext().getExecutionVariables()); - executionVariableDescriptors.clear(); - executionVariableDescriptors.putAll(innerExecution.getContext().getExecutionVariableDescriptors()); - - for (IteratorContainerInterfaceResolver.ResolvedOutput resolvedOutput : resolution.resolvedOutputs()) { - Object value = innerExecution.getContext().getResult() - .get(new FieldKey( - resolvedOutput.exposedHandle().handle().blockId(), - resolvedOutput.exposedHandle().handle().io().getName())); - collectedOutputs.get(resolvedOutput.publicName()).add(value); - } - if (eventLogger != null) { - eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, - "Completed iterator iteration " + iterationIndex, - Map.of("containerType", IteratorContainerType.TYPE, "iteration", iterationIndex)); + Map runtimeInputValues = new LinkedHashMap<>(); + for (IteratorContainerInterfaceResolver.ResolvedInput resolvedInput : resolution.resolvedInputs()) { + if (!resolvedInput.iterated()) { + runtimeInputValues.put(resolvedInput.publicName(), inputs.stream() + .filter(input -> input.getDescriptor().getName().equals(resolvedInput.publicName())) + .findFirst() + .map(Input::getValue) + .orElse(null)); } } - return NodeExecutionResult.completed(new LinkedHashMap<>(collectedOutputs)); - } - - private void propagateAuthorizations(ExecutionObject innerExecution, Map authorizations) { - for (String key : innerExecution.getRequiredAuthorizations().stream().map(requirement -> requirement.key()).distinct() - .toList()) { - if (authorizations.containsKey(key)) { - executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key)); - } - } - } - - private ExecutionObject startAndWait(ExecutionObject innerExecution) { - long startedAtNanos = System.nanoTime(); - String finalStatus = "unknown"; - innerExecution = executionsService.startExecution(innerExecution.getId()); - try { - while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) { - innerExecution.getContext().awaitStatusChangeWhileRunning(innerExecutionWaitIntervalMs); - } - - if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) { - finalStatus = "waiting"; - throw new IllegalStateException( - "IteratorContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet"); - } - if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) { - finalStatus = "error"; - throw new IllegalStateException("IteratorContainer subflow failed: " + innerExecution.getContext().getErrors()); - } - if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) { - finalStatus = innerExecution.getContext().getStatus().name().toLowerCase(); - throw new IllegalStateException( - "IteratorContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus()); - } - finalStatus = "success"; - return innerExecution; - } finally { - long elapsedNanos = System.nanoTime() - startedAtNanos; - recordSubflowDuration(finalStatus, elapsedNanos); - long elapsedMs = elapsedNanos / 1_000_000L; - if (elapsedMs >= SLOW_SUBFLOW_THRESHOLD_MS) { - logger.warn("Iterator container subflow was slow: executionId={}, status={}, elapsedMs={}", - innerExecution.getId(), finalStatus, elapsedMs); - } - } - } - - private void recordSubflowDuration(String status, long elapsedNanos) { - if (meterRegistry == null) { - return; - } - Timer.builder("workflow.container.subflow.duration") - .description("Duration of container subflow executions") - .tag("containerType", IteratorContainerType.TYPE) - .tag("status", status) - .register(meterRegistry) - .record(elapsedNanos, java.util.concurrent.TimeUnit.NANOSECONDS); - } - - private void forwardInnerEvents(ExecutionObject innerExecution, ExecutionEventLogger eventLogger, Integer iterationIndex, - String containerType) { - if (eventLogger == null) { - return; - } - innerExecution.getContext().addEventListener(event -> eventLogger.log( - event.getLevel(), - event.getType(), - prefixMessage(event), - enrichDetails(event, innerExecution.getId(), iterationIndex, containerType))); - } - - private String prefixMessage(ExecutionEvent event) { - return event.getNodeName() == null || event.getNodeName().isBlank() - ? event.getMessage() - : "[" + event.getNodeName() + "] " + event.getMessage(); - } - - private Map enrichDetails(ExecutionEvent event, String innerExecutionId, Integer iterationIndex, - String containerType) { - Map details = new LinkedHashMap<>(); - details.put("containerType", containerType); - details.put("innerExecutionId", innerExecutionId); - details.put("iterationIndex", iterationIndex); - if (event.getNodeId() != null) { - details.put("innerNodeId", event.getNodeId()); - } - if (event.getNodeName() != null) { - details.put("innerNodeName", event.getNodeName()); - } - if (event.getStepId() != null) { - details.put("innerStepId", event.getStepId()); - } - if (event.getDetails() != null && !event.getDetails().isEmpty()) { - details.putAll(event.getDetails()); - } - return details; + return executionsService.startIteratorSubflow(executionContext.parentExecutionId(), executionContext.parentStepId(), + container, configuration, resolution, new java.util.ArrayList<>(iterationValues), runtimeInputValues, + authorizations, eventLogger, innerBiasExecutionContext, executionContext); } @Override public Class getContainerType() { return IteratorContainerType.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; - }; - } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java index 9c4039b..bcf0319 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/containers/LoopContainerExecutor.java @@ -3,52 +3,35 @@ package it.cnr.isti.workflow.manager.executions.executors.containers; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; -import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.Timer; - import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration; -import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; -import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ContainerExecutionContext; -import it.cnr.isti.workflow.manager.executions.NodeExecutionResult; -import it.cnr.isti.workflow.manager.executions.ExecutionEvent; -import it.cnr.isti.workflow.manager.executions.ExecutionEventType; -import it.cnr.isti.workflow.manager.executions.ExecutionObject; -import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport; -import it.cnr.isti.workflow.manager.executions.ExecutionStatus; +import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; import it.cnr.isti.workflow.manager.executions.ExecutionsService; -import it.cnr.isti.workflow.manager.executions.FieldKey; +import it.cnr.isti.workflow.manager.executions.NodeExecutionResult; 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.executors.BooleanLlmResponseParser; import it.cnr.isti.workflow.manager.executions.steps.Input; +/** + * Runs the main subflow followed by the guard subflow on each iteration. + * Iterations that complete synchronously (no interactive node reached) are + * chained without suspending the step; the moment one doesn't, the step + * suspends and is resumed later by {@link ExecutionsService}'s durable + * container coordinator. + */ @Component public class LoopContainerExecutor implements ContainerExecutor { - private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(LoopContainerExecutor.class); - private static final long SLOW_SUBFLOW_THRESHOLD_MS = 3_000L; - private final ExecutionsService executionsService; - private final long innerExecutionWaitIntervalMs; - @Autowired(required = false) - MeterRegistry meterRegistry; - - public LoopContainerExecutor( - ExecutionsService executionsService, - @Value("${app.container.subflow.wait-interval-ms:${app.container.subflow.wait-timeout-ms:5000}}") long innerExecutionWaitIntervalMs) { + public LoopContainerExecutor(ExecutionsService executionsService) { this.executionsService = executionsService; - this.innerExecutionWaitIntervalMs = innerExecutionWaitIntervalMs; } @Override @@ -66,339 +49,18 @@ public class LoopContainerExecutor implements ContainerExecutor inputPortsByName = ContainerFlowInterfaceResolver - .getExposedInputs(configuration.getSubFlow()).stream() - .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed)); - String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName); - List exposedOutputs = ContainerFlowInterfaceResolver - .getExposedOutputs(configuration.getSubFlow()); - - Map currentInputs = inputs.stream() + Map initialInputs = inputs.stream() .collect(LinkedHashMap::new, (map, input) -> map.put(input.getDescriptor().getName(), input.getValue()), Map::putAll); - Map latestOutputs = new LinkedHashMap<>(); - for (int iteration = 1; iteration <= configuration.getMaxIterations(); iteration++) { - if (eventLogger != null) { - eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_STARTED, - "Starting loop iteration " + iteration, - Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration)); - } - - ExecutionObject innerExecution = executeSubFlow( - container.getName() + " iteration " + iteration, - configuration.getSubFlow(), - currentInputs, - inputPortsByName, - "LoopContainer", - authorizations, - executionVariables, - executionVariableDescriptors, - eventLogger, - iteration, - innerBiasExecutionContext); - - latestOutputs.clear(); - for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedOutputs) { - ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle(); - Object value = innerExecution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); - if (value != null) { - latestOutputs.put(exposedHandle.publicName(), value); - } - } - - Map iterationInputs = new LinkedHashMap<>(currentInputs); - updateInputsForNextIteration(currentInputs, latestOutputs, inputPortsByName); - - GuardDecision decision = evaluateGuardSubFlow(configuration, container.getName(), iterationInputs, latestOutputs, - authorizations, executionVariables, executionVariableDescriptors, eventLogger, iteration, - innerBiasExecutionContext); - if (eventLogger != null) { - eventLogger.info(ExecutionEventType.CONTAINER_CONDITION_EVALUATED, - "Evaluated loop guard", - Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "continue", decision.shouldContinue())); - eventLogger.info(ExecutionEventType.CONTAINER_ITERATION_COMPLETED, - decision.shouldContinue() - ? "Completed loop iteration " + iteration - : "Loop completed at iteration " + iteration, - Map.of("containerType", LoopContainerType.TYPE, "iteration", iteration, "continue", - decision.shouldContinue())); - } - if (decision.shouldContinue()) { - currentInputs.put(feedbackInput, decision.feedback()); - } else { - return NodeExecutionResult.completed(new LinkedHashMap<>(latestOutputs)); - } - } - - throw new IllegalStateException( - "LoopContainer guard was not satisfied within maxIterations=" + configuration.getMaxIterations()); + return executionsService.startLoopSubflow(executionContext.parentExecutionId(), executionContext.parentStepId(), + container, configuration, initialInputs, authorizations, eventLogger, innerBiasExecutionContext, + executionContext); } @Override public Class getContainerType() { return LoopContainerType.class; } - - private void updateInputsForNextIteration(Map currentInputs, Map latestOutputs, - Map inputPortsByName) { - for (Map.Entry output : latestOutputs.entrySet()) { - if (inputPortsByName.containsKey(output.getKey())) { - currentInputs.put(output.getKey(), output.getValue()); - } - } - } - - private String resolveFeedbackInput(LoopContainerConfiguration configuration, - Map inputPortsByName) { - if (configuration.getFeedbackInput() != null && !configuration.getFeedbackInput().isBlank()) { - String configuredInput = configuration.getFeedbackInput(); - ContainerFlowInterfaceResolver.ExposedHandle handle = inputPortsByName.get(configuredInput); - if (handle == null || handle.handle().io().isMultiple()) { - throw new IllegalArgumentException( - "LoopContainer feedbackInput must target an open non-multiple subFlow input: " + configuredInput); - } - return configuredInput; - } - - List eligibleInputs = inputPortsByName.values().stream() - .filter(handle -> !handle.handle().io().isMultiple()) - .toList(); - if (eligibleInputs.size() == 1) { - return eligibleInputs.getFirst().publicName(); - } - if (eligibleInputs.isEmpty()) { - throw new IllegalArgumentException( - "LoopContainer subFlow must expose at least one open non-multiple input to receive guard feedback"); - } - throw new IllegalArgumentException( - "LoopContainer feedbackInput is required when subFlow exposes more than one open non-multiple input"); - } - - private GuardDecision evaluateGuardSubFlow(LoopContainerConfiguration configuration, String containerName, - Map iterationInputs, Map latestOutputs, Map authorizations, - Map executionVariables, - Map executionVariableDescriptors, ExecutionEventLogger eventLogger, - int iteration, BiasExecutionContext innerBiasExecutionContext) { - Map guardInputsByName = ContainerFlowInterfaceResolver - .getExposedInputs(configuration.getGuardSubFlow()).stream() - .collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed)); - List guardOutputs = ContainerFlowInterfaceResolver - .getExposedOutputs(configuration.getGuardSubFlow()); - Map guardInputValues = buildGuardTemplateValues(iterationInputs, latestOutputs, iteration); - ExecutionObject guardExecution = executeSubFlow( - containerName + " guard iteration " + iteration, - configuration.getGuardSubFlow(), - guardInputValues, - guardInputsByName, - "LoopContainerGuard", - authorizations, - executionVariables, - executionVariableDescriptors, - eventLogger, - iteration, - innerBiasExecutionContext); - Map guardResult = collectExposedOutputs(guardExecution, guardOutputs); - Object guard = requireGuardOutput(guardResult, LoopContainerConfiguration.GUARD_OUTPUT); - Object feedback = requireGuardOutput(guardResult, LoopContainerConfiguration.FEEDBACK_OUTPUT); - return new GuardDecision(parseBooleanResponse(String.valueOf(guard)), String.valueOf(feedback)); - } - - private Map buildGuardTemplateValues(Map currentInputs, - Map latestOutputs, int iteration) { - Map values = new LinkedHashMap<>(); - if (currentInputs != null) { - values.putAll(currentInputs); - currentInputs.forEach((key, value) -> values.put("inputs." + key, value)); - } - if (latestOutputs != null) { - values.putAll(latestOutputs); - latestOutputs.forEach((key, value) -> values.put("outputs." + key, value)); - } - values.put("iteration", iteration); - return values; - } - - private boolean parseBooleanResponse(String response) { - return BooleanLlmResponseParser.parse(response, "loop guard subflow", "Loop"); - } - - private ExecutionObject executeSubFlow(String executionName, it.cnr.isti.workflow.manager.flows.model.FlowData subFlow, - Map inputValues, - Map inputPortsByName, - String containerType, - Map authorizations, - Map executionVariables, - Map executionVariableDescriptors, - ExecutionEventLogger eventLogger, - int iteration, - BiasExecutionContext innerBiasExecutionContext) { - ExecutionObject innerExecution = executionsService.createInnerExecution(executionName, subFlow, innerBiasExecutionContext); - forwardInnerEvents(innerExecution, eventLogger, iteration, containerType); - - propagateAuthorizations(innerExecution, authorizations); - executionsService.setExecutionVariableDescriptors(innerExecution.getId(), executionVariableDescriptors); - Map globalInputs = ExecutionRuntimeContextSupport.globalView(executionVariables); - executionsService.setGlobalInputDescriptors(innerExecution.getId(), - innerExecution.getRequiredGlobalInputs().stream() - .collect(LinkedHashMap::new, (map, required) -> map.put(required.getName(), - ExecutionVariableDescriptor.builder() - .name(required.getName()) - .kind(mapKind(required)) - .value(globalInputs.get(required.getName())) - .cleanupPolicy(it.cnr.isti.workflow.manager.executions.ExecutionVariableCleanupPolicy.NONE) - .description("Global flow input") - .build()), Map::putAll)); - - for (Map.Entry entry : inputPortsByName.entrySet()) { - if (!inputValues.containsKey(entry.getKey())) { - throw new IllegalArgumentException( - containerType + " subflow input is missing: " + entry.getKey()); - } - ContainerFlowInterfaceResolver.OpenHandle handle = entry.getValue().handle(); - executionsService.prepareInput(innerExecution.getId(), handle.blockId(), handle.io().getName(), - inputValues.get(entry.getKey())); - } - - innerExecution = startAndWait(innerExecution); - executionVariables.clear(); - executionVariables.putAll(innerExecution.getContext().getExecutionVariables()); - executionVariableDescriptors.clear(); - executionVariableDescriptors.putAll(innerExecution.getContext().getExecutionVariableDescriptors()); - return innerExecution; - } - - private Map collectExposedOutputs(ExecutionObject execution, - List exposedOutputs) { - Map outputs = new LinkedHashMap<>(); - for (ContainerFlowInterfaceResolver.ExposedHandle exposedHandle : exposedOutputs) { - ContainerFlowInterfaceResolver.OpenHandle handle = exposedHandle.handle(); - Object value = execution.getContext().getResult().get(new FieldKey(handle.blockId(), handle.io().getName())); - if (value != null) { - outputs.put(exposedHandle.publicName(), value); - } - } - return outputs; - } - - private Object requireGuardOutput(Map guardResult, String outputName) { - Object value = guardResult.get(outputName); - if (value == null) { - throw new IllegalArgumentException("LoopContainer guardSubFlow must return output: " + outputName); - } - return value; - } - - private void propagateAuthorizations(ExecutionObject innerExecution, Map authorizations) { - for (String key : innerExecution.getRequiredAuthorizations().stream() - .map(requirement -> requirement.key()) - .distinct() - .toList()) { - if (authorizations.containsKey(key)) { - executionsService.setAuthorizationValue(innerExecution.getId(), key, authorizations.get(key)); - } - } - } - - private ExecutionObject startAndWait(ExecutionObject innerExecution) { - long startedAtNanos = System.nanoTime(); - String finalStatus = "unknown"; - innerExecution = executionsService.startExecution(innerExecution.getId()); - try { - while (innerExecution.getContext().getStatus() == ExecutionStatus.RUNNING) { - innerExecution.getContext().awaitStatusChangeWhileRunning(innerExecutionWaitIntervalMs); - } - - if (innerExecution.getContext().getStatus() == ExecutionStatus.WAITING) { - finalStatus = "waiting"; - throw new IllegalStateException( - "LoopContainer subflow reached WAITING state. Interactive blocks inside containers are not supported yet"); - } - if (innerExecution.getContext().getStatus() == ExecutionStatus.ERROR) { - finalStatus = "error"; - throw new IllegalStateException("LoopContainer subflow failed: " + innerExecution.getContext().getErrors()); - } - if (innerExecution.getContext().getStatus() != ExecutionStatus.SUCCESS) { - finalStatus = innerExecution.getContext().getStatus().name().toLowerCase(); - throw new IllegalStateException( - "LoopContainer subflow ended in unexpected status: " + innerExecution.getContext().getStatus()); - } - finalStatus = "success"; - return innerExecution; - } finally { - long elapsedNanos = System.nanoTime() - startedAtNanos; - recordSubflowDuration(finalStatus, elapsedNanos); - long elapsedMs = elapsedNanos / 1_000_000L; - if (elapsedMs >= SLOW_SUBFLOW_THRESHOLD_MS) { - logger.warn("Loop container subflow was slow: executionId={}, status={}, elapsedMs={}", - innerExecution.getId(), finalStatus, elapsedMs); - } - } - } - - private void recordSubflowDuration(String status, long elapsedNanos) { - if (meterRegistry == null) { - return; - } - Timer.builder("workflow.container.subflow.duration") - .description("Duration of container subflow executions") - .tag("containerType", LoopContainerType.TYPE) - .tag("status", status) - .register(meterRegistry) - .record(elapsedNanos, java.util.concurrent.TimeUnit.NANOSECONDS); - } - - private void forwardInnerEvents(ExecutionObject innerExecution, ExecutionEventLogger eventLogger, Integer iterationIndex, - String containerType) { - if (eventLogger == null) { - return; - } - innerExecution.getContext().addEventListener(event -> eventLogger.log( - event.getLevel(), - event.getType(), - prefixMessage(event), - enrichDetails(event, innerExecution.getId(), iterationIndex, containerType))); - } - - 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; - }; - } - - private String prefixMessage(ExecutionEvent event) { - return event.getNodeName() == null || event.getNodeName().isBlank() - ? event.getMessage() - : "[" + event.getNodeName() + "] " + event.getMessage(); - } - - private Map enrichDetails(ExecutionEvent event, String innerExecutionId, Integer iterationIndex, - String containerType) { - Map details = new LinkedHashMap<>(); - details.put("containerType", containerType); - details.put("innerExecutionId", innerExecutionId); - details.put("iterationIndex", iterationIndex); - if (event.getNodeId() != null) { - details.put("innerNodeId", event.getNodeId()); - } - if (event.getNodeName() != null) { - details.put("innerNodeName", event.getNodeName()); - } - if (event.getStepId() != null) { - details.put("innerStepId", event.getStepId()); - } - if (event.getDetails() != null && !event.getDetails().isEmpty()) { - details.putAll(event.getDetails()); - } - return details; - } - - private record GuardDecision(boolean shouldContinue, String feedback) { - } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java index ce2ac03..c8a7478 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java @@ -100,6 +100,17 @@ public class Step implements InputListener { @Getter private boolean simulated = false; + /** + * Whether the containing execution overall is running in simulation + * mode. Unlike {@link #simulated}, this is set on every step (containers + * included) so a container can propagate simulation to its inner + * subflow even though the container node itself is never + * {@code isUserInteractive()}. + */ + @Setter + @Getter + private boolean executionSimulationEnabled = false; + @Setter @Getter @JsonIgnore @@ -126,6 +137,7 @@ public class Step implements InputListener { private BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal(); @Setter + @Getter @JsonIgnore private ExecutionEventLogger eventLogger; @@ -209,7 +221,8 @@ public class Step implements InputListener { this.biasExecutionContext) : NodeExecutors.executeResult(this.node, this.inputs, authorizations, executionVariables, executionVariableDescriptors, this.eventLogger, this.biasExecutionContext, - new ContainerExecutionContext(this.parentExecutionId, this.id)); + new ContainerExecutionContext(this.parentExecutionId, this.id, + this.executionSimulationEnabled, this.interactionSimulationDescriptor)); if (executionResult.suspended()) { this.status = StepStatus.WAITING_FOR_SUBFLOW; listener.paused(this.id); diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java index 3af277c..e42da47 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java @@ -25,6 +25,7 @@ import it.cnr.isti.workflow.manager.blocks.types.SwitchBlockType; import it.cnr.isti.workflow.manager.containers.Container; import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.GenericContainerConfiguration; +import it.cnr.isti.workflow.manager.containers.configurations.IteratorContainerConfiguration; import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration; import it.cnr.isti.workflow.manager.containers.factories.ContainerFactory; import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver; @@ -397,8 +398,14 @@ public class FlowDataValidator implements ConstraintValidator}); the default + * {@code loopGuardSubFlow()} assumes "response" + * (the LLMBlockFactory output name). + */ + private FlowData loopGuardSubFlow(String mainOutputName) { Block guardEvaluator = llmBlockFactory.create(LLMBlockConfiguration.builder() .name("Guard Evaluator") .llmDescriptor(llmBrick) - .prompt("Loop guard subflow. Previous output: ${{outputs.response}}") + .prompt("Loop guard subflow. Previous output: ${{outputs." + mainOutputName + "}}") .build()); Block guardOutput = switchBlockFactory.create(SwitchBlockConfiguration.builder() .name("Expose Guard") @@ -1412,6 +1425,209 @@ public class ExecutionTest { return execution; } + @Test + public void loopContainerSuspendsForInnerInteractionAndResumesFromChild() { + Block innerInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("Inner review") + .actionDescription("Review the submitted value") + .build()); + Container container = loopContainerFactory.create(LoopContainerConfiguration.builder() + .name("Interactive loop") + .subFlow(FlowData.builder().block(innerInteraction).build()) + .guardSubFlow(loopGuardSubFlow("output")) + .maxIterations(3) + .build()); + ExecutionObject parent = executionsService.createExecution("Interactive loop flow", + FlowData.builder().container(container).build(), "loop-owner"); + + executionsService.prepareInput(parent.getId(), container.getId(), "input", "submission"); + executionsService.startExecution(parent.getId()); + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.WAITING); + Step containerStep = parent.getContext().getSteps().get(container.getId()); + assertEquals(StepStatus.WAITING_FOR_SUBFLOW, containerStep.getStatus()); + assertNotNull(containerStep.getContainerContinuation()); + String childId = containerStep.getContainerContinuation().getActiveInnerExecutionId(); + + ExecutionObject child = executionsService.getExecutionByOwner(childId, "loop-owner"); + assertEquals(ExecutionKind.SUBFLOW, child.getExecutionKind()); + assertEquals(ContainerSubflowRole.MAIN, child.getSubflowRole()); + assertEquals(StepStatus.WAITING_FOR_INTERACTION, + child.getContext().getSteps().get(innerInteraction.getId()).getStatus()); + + // The guard mock (see loopGuardSubFlow()) returns false unless the main + // output contains "function v1", so a single "approved" iteration stops + // the loop without needing a second (interactive) iteration. + executionsService.setInteractionValue(childId, innerInteraction.getId(), "output", "approved"); + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.SUCCESS); + + assertEquals("approved", parent.getContext().getResult().get(new FieldKey(container.getId(), "output"))); + assertTrue(parent.getContext().getEvents().stream() + .anyMatch(event -> event.getType() == ExecutionEventType.STEP_WAITING_FOR_SUBFLOW)); + } + + @Test + public void iteratorContainerSuspendsPerIterationForInnerInteractionAndAccumulatesOutputs() { + Block innerInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("Inner review") + .actionDescription("Review the submitted value") + .build()); + Container container = iteratorContainerFactory.create(IteratorContainerConfiguration.builder() + .name("Interactive iterator") + .subFlow(FlowData.builder().block(innerInteraction).build()) + .iterationInput("input") + .build()); + ExecutionObject parent = executionsService.createExecution("Interactive iterator flow", + FlowData.builder().container(container).build(), "iterator-owner"); + + executionsService.prepareInput(parent.getId(), container.getId(), "input", List.of("first", "second")); + executionsService.startExecution(parent.getId()); + + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.WAITING); + Step containerStep = parent.getContext().getSteps().get(container.getId()); + assertEquals(StepStatus.WAITING_FOR_SUBFLOW, containerStep.getStatus()); + assertEquals(1, containerStep.getContainerContinuation().getIterationIndex()); + String firstChildId = containerStep.getContainerContinuation().getActiveInnerExecutionId(); + executionsService.setInteractionValue(firstChildId, innerInteraction.getId(), "output", "approved-1"); + + parent = executionsService.getExecution(parent.getId()); + containerStep = parent.getContext().getSteps().get(container.getId()); + long deadline = System.currentTimeMillis() + 10_000L; + while (containerStep.getContainerContinuation() != null + && firstChildId.equals(containerStep.getContainerContinuation().getActiveInnerExecutionId()) + && System.currentTimeMillis() < deadline) { + try { + Thread.sleep(10L); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new RuntimeException(exception); + } + parent = executionsService.getExecution(parent.getId()); + containerStep = parent.getContext().getSteps().get(container.getId()); + } + assertEquals(StepStatus.WAITING_FOR_SUBFLOW, containerStep.getStatus()); + assertEquals(2, containerStep.getContainerContinuation().getIterationIndex()); + String secondChildId = containerStep.getContainerContinuation().getActiveInnerExecutionId(); + assertNotEquals(firstChildId, secondChildId); + executionsService.setInteractionValue(secondChildId, innerInteraction.getId(), "output", "approved-2"); + + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.SUCCESS); + Object output = parent.getContext().getResult().get(new FieldKey(container.getId(), "output")); + assertEquals(List.of("approved-1", "approved-2"), output); + } + + @Test + public void genericContainerPropagatesSimulationAndDoesNotSuspendForInnerInteraction() { + Block innerInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("Inner review") + .actionDescription("Review the submitted value") + .build()); + Container container = genericContainerFactory.create(GenericContainerConfiguration.builder() + .name("Simulated container") + .subFlow(FlowData.builder().block(innerInteraction).build()) + .build()); + ExecutionObject parent = executionsService.createExecution("Simulated container flow", + FlowData.builder().container(container).build(), "simulation-owner"); + + executionsService.prepareInput(parent.getId(), container.getId(), "input", "submission"); + assertTrue(parent.isSimulationAvailable()); + executionsService.startSimulationExecution(parent.getId(), SIMULATOR_DESCRIPTOR); + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.SUCCESS); + + Step containerStep = parent.getContext().getSteps().get(container.getId()); + assertEquals(StepStatus.COMPLETED, containerStep.getStatus()); + assertNull(containerStep.getContainerContinuation()); + } + + @Test + public void cancellingParentCancelsActiveChildAndCancellingChildFailsParentStep() { + Block firstInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("First review") + .actionDescription("Review the submitted value") + .build()); + Container firstContainer = genericContainerFactory.create(GenericContainerConfiguration.builder() + .name("Cancellable container") + .subFlow(FlowData.builder().block(firstInteraction).build()) + .build()); + ExecutionObject parent = executionsService.createExecution("Parent cancellation flow", + FlowData.builder().container(firstContainer).build(), "cancel-owner"); + executionsService.prepareInput(parent.getId(), firstContainer.getId(), "input", "submission"); + executionsService.startExecution(parent.getId()); + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.WAITING); + String childId = parent.getContext().getSteps().get(firstContainer.getId()) + .getContainerContinuation().getActiveInnerExecutionId(); + + executionsService.cancelExecution(parent.getId()); + + ExecutionObject cancelledChild = executionsService.getExecutionByOwner(childId, "cancel-owner"); + assertEquals(ExecutionStatus.CANCELLED, cancelledChild.getContext().getStatus()); + assertEquals(ExecutionStatus.CANCELLED, executionsService.getExecution(parent.getId()).getContext().getStatus()); + + // Symmetric direction: cancelling the child directly fails the parent's + // container step (a cancelled subflow is not a successful completion). + Block secondInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("Second review") + .actionDescription("Review the submitted value") + .build()); + Container secondContainer = genericContainerFactory.create(GenericContainerConfiguration.builder() + .name("Cancellable container 2") + .subFlow(FlowData.builder().block(secondInteraction).build()) + .build()); + ExecutionObject secondParent = executionsService.createExecution("Child cancellation flow", + FlowData.builder().container(secondContainer).build(), "cancel-owner"); + executionsService.prepareInput(secondParent.getId(), secondContainer.getId(), "input", "submission"); + executionsService.startExecution(secondParent.getId()); + secondParent = waitForExecutionStatus(secondParent.getId(), ExecutionStatus.WAITING); + String secondChildId = secondParent.getContext().getSteps().get(secondContainer.getId()) + .getContainerContinuation().getActiveInnerExecutionId(); + + executionsService.cancelExecution(secondChildId); + + secondParent = waitForExecutionStatus(secondParent.getId(), ExecutionStatus.ERROR); + assertTrue(secondParent.getContext().getErrors().containsKey(secondContainer.getId())); + } + + @Test + public void containerCoordinatorReconcilesChildCompletionAfterCacheEviction() { + Block innerInteraction = humanInteractiveBlockFactory.create( + HumanInteractiveBlockConfiguration.builder() + .name("Inner review") + .actionDescription("Review the submitted value") + .build()); + Container container = genericContainerFactory.create(GenericContainerConfiguration.builder() + .name("Durable container") + .subFlow(FlowData.builder().block(innerInteraction).build()) + .build()); + ExecutionObject parent = executionsService.createExecution("Durable reconciliation flow", + FlowData.builder().container(container).build(), "durable-owner"); + executionsService.prepareInput(parent.getId(), container.getId(), "input", "submission"); + executionsService.startExecution(parent.getId()); + parent = waitForExecutionStatus(parent.getId(), ExecutionStatus.WAITING); + String childId = parent.getContext().getSteps().get(container.getId()) + .getContainerContinuation().getActiveInnerExecutionId(); + + // Simulate a restart: drop every in-memory execution (and any listener + // attached to it) so nothing but the persisted parent-child link and + // continuation remain, then run the same reconciliation the real + // ApplicationReadyEvent listener runs at startup. + executionsService.clearInMemoryExecutions(); + executionsService.reconcileContainerSubflowsOnStartup(); + + ExecutionObject reloadedChild = executionsService.getExecutionByOwner(childId, "durable-owner"); + assertEquals(StepStatus.WAITING_FOR_INTERACTION, + reloadedChild.getContext().getSteps().get(innerInteraction.getId()).getStatus()); + executionsService.resumeExecution(childId); + executionsService.setInteractionValue(childId, innerInteraction.getId(), "output", "approved-after-restart"); + + ExecutionObject reloadedParent = waitForExecutionStatus(parent.getId(), ExecutionStatus.SUCCESS); + assertEquals("approved-after-restart", + reloadedParent.getContext().getResult().get(new FieldKey(container.getId(), "output"))); + } + @Test public void iteratorContainerExecutesSubflowForEachIteratedValue() { Block internalBlock = llmBlockFactory.create(LLMBlockConfiguration.builder() @@ -1516,9 +1732,7 @@ public class ExecutionTest { executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Alice"); execObject = executionsService.startExecution(execObject.getId()); - while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { - execObject = executionsService.getExecution(execObject.getId()); - } + execObject = awaitNonBlockingContainerCompletion(execObject); assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); @@ -1545,9 +1759,7 @@ public class ExecutionTest { executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Alice"); execObject = executionsService.startExecution(execObject.getId()); - while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { - execObject = executionsService.getExecution(execObject.getId()); - } + execObject = awaitNonBlockingContainerCompletion(execObject); assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); @@ -1574,9 +1786,7 @@ public class ExecutionTest { executionsService.prepareInput(execObject.getId(), container.getId(), "feedback", "Implement the same function"); execObject = executionsService.startExecution(execObject.getId()); - while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { - execObject = executionsService.getExecution(execObject.getId()); - } + execObject = awaitNonBlockingContainerCompletion(execObject); assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); @@ -1607,9 +1817,7 @@ public class ExecutionTest { executionsService.prepareInput(execObject.getId(), container.getId(), "specification", "Implement the same function"); executionsService.prepareInput(execObject.getId(), container.getId(), "language", "Java"); execObject = executionsService.startExecution(execObject.getId()); - while (execObject.getContext().getStatus() == ExecutionStatus.RUNNING) { - execObject = executionsService.getExecution(execObject.getId()); - } + execObject = awaitNonBlockingContainerCompletion(execObject); assertEquals(ExecutionStatus.SUCCESS, execObject.getContext().getStatus()); Object output = execObject.getContext().getResult().values().stream().findFirst().orElseThrow(); @@ -1630,6 +1838,29 @@ public class ExecutionTest { assertTrue(exception.getMessage().contains("does not accept multiple values")); } + /** + * A container's subflow always runs on its own thread pool, so even a + * fully automatic container (no real interactive node) transiently visits + * WAITING while its child executes, resolved shortly after by the + * container's durable reconciliation listener. Polls through both RUNNING + * and WAITING, bounded so a genuine hang still fails fast. + */ + private ExecutionObject awaitNonBlockingContainerCompletion(ExecutionObject execObject) { + long deadline = System.currentTimeMillis() + 30_000; + while ((execObject.getContext().getStatus() == ExecutionStatus.RUNNING + || execObject.getContext().getStatus() == ExecutionStatus.WAITING) + && System.currentTimeMillis() < deadline) { + try { + Thread.sleep(20); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + execObject = executionsService.getExecution(execObject.getId()); + } + return execObject; + } + private ExecutionObject createExecutionAndSetInputInternally() { Flow flow = flowTestCreator.createFlowWithConnection(llmBrick); ExecutionObject execObject = executionsService.createExecution(flow);