feat: durable container coordinator and non-blocking Loop/Iterator subflows
Closes Tappa B and Tappa C of the interactive-containers plan. Tappa B (recovery and lifecycle): - Generalize the GenericContainer-only watcher into a type-dispatching coordinator (watchContainerSubflow/reconcileContainerSubflow) shared by Generic, Loop and Iterator. - Add reconcileContainerSubflowsOnStartup (ApplicationReadyEvent): re-arms every persisted parent-child link at boot using only the entity's own parentExecutionId/parentStepId columns, so a child that already finished while the process was down is reconciled immediately, and a pending one gets its listener re-armed. No longer relies on any in-memory listener surviving a restart. - cancelExecution now propagates in both directions: cancelling a parent cancels its active container child; cancelling a child fails the parent's container step (same outcome as an unexpected child error). removeExecution evicts cached children too. - Propagate simulation mode to container children: Step gains executionSimulationEnabled (mirrors interactionSimulationDescriptor's propagation to every step, containers included); ContainerExecutionContext carries it plus the descriptor; startContainerChild starts the child via startSimulationExecution when applicable. ExecutionObject.hasSimulationAvailable now recognizes interactive nodes nested inside a container's subflow(s), which it previously ignored entirely since a container is never itself isUserInteractive(). Tappa C (Loop/Iterator non-blocking): - LoopContainerExecutor and IteratorContainerExecutor become thin adapters delegating to ExecutionsService.startLoopSubflow/startIteratorSubflow. All advancement logic (creating/starting a child, checking whether it finished synchronously vs. suspending) now lives in ExecutionsService, so it is reusable both from the initial call and from reconciliation. - Iterator advances iteration-by-iteration in a loop, chaining synchronous completions without suspending, and persists remainingValues/ runtimeInputValues/accumulatedOutputs so it can resume exactly where it left off. - Loop is a two-phase (MAIN/GUARD) state machine; the phase to resume into is derived from the completed child's own subflowRole, so it doesn't need separate persistence. Persists currentInputs/latestOutputs. - FlowDataValidator and FlowExecutionValidator now accept interactive nodes in the subFlow of all three container types; LoopContainer's guardSubFlow remains rejected (guard interactivity is out of scope). Existing container tests that assumed the old fully-synchronous model (assert SUCCESS right after start()) are updated: since a child's steps always run on their own thread pool, even a fully-automatic container now transiently visits WAITING before the coordinator resolves it, so tests must poll through WAITING too, not just RUNNING. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
parent
0e5831289f
commit
a31e1692f6
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<Step<?>> 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<FlowData> 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();
|
||||
|
|
|
|||
|
|
@ -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<String> 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<String> 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<String> 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<String, Object> 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<Object> iterationValues,
|
||||
Map<String, Object> runtimeInputValues, Map<String, Object> authorizations, ExecutionEventLogger eventLogger,
|
||||
BiasExecutionContext innerBiasExecutionContext, ContainerExecutionContext executionContext) {
|
||||
Map<String, List<Object>> 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<Object> 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<ContainerFlowInterfaceResolver.ExposedHandle> outputHandles = resolution.resolvedOutputs().stream()
|
||||
.map(IteratorContainerInterfaceResolver.ResolvedOutput::exposedHandle).toList();
|
||||
Map<String, Object> iterationOutputs = collectExposedOutputsAsMap(child, outputHandles);
|
||||
Map<String, List<Object>> 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<String, Object> state = continuation.getState() == null ? Map.of() : continuation.getState();
|
||||
List<Object> remaining = asObjectList(state.get("remainingValues"));
|
||||
Map<String, Object> 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<Object> 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<Object> remainingValues, Map<String, Object> runtimeInputValues, Map<String, List<Object>> accumulatedOutputs,
|
||||
Map<String, Object> authorizations, ExecutionEventLogger eventLogger, BiasExecutionContext innerBiasExecutionContext,
|
||||
ContainerExecutionContext executionContext) {
|
||||
int iterationIndex = startIterationIndex;
|
||||
Object iterationValue = startValue;
|
||||
List<Object> remaining = remainingValues;
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = resolution.resolvedInputs().stream()
|
||||
.collect(Collectors.toMap(IteratorContainerInterfaceResolver.ResolvedInput::publicName,
|
||||
IteratorContainerInterfaceResolver.ResolvedInput::exposedHandle));
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> 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<String, Object> 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<String, Object> initialInputs, Map<String, Object> 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<String, Object> state = continuation.getState() == null ? Map.of() : continuation.getState();
|
||||
Map<String, Object> currentInputs = asStringObjectMap(state.get("currentInputs"));
|
||||
Map<String, Object> 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<String, Object> authorizations = parent.getProvidedAuthorizations();
|
||||
|
||||
ContainerAdvanceOutcome outcome;
|
||||
if (completedChild.getSubflowRole() != ContainerSubflowRole.GUARD) {
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> mainOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getSubFlow());
|
||||
Map<String, Object> freshLatest = collectExposedOutputsAsMap(completedChild, mainOutputHandles);
|
||||
outcome = runLoopFrom(parent, parentStep.getId(), container, configuration, LoopPhase.GUARD, iterationIndex,
|
||||
currentInputs, freshLatest, authorizations, eventLogger, innerBiasExecutionContext, executionContext);
|
||||
} else {
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getSubFlow()).stream()
|
||||
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle));
|
||||
String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> guardOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getGuardSubFlow());
|
||||
Map<String, Object> 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<String, Object> 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<String, Object> startInputs, Map<String, Object> startLatestOutputs, Map<String, Object> authorizations,
|
||||
ExecutionEventLogger eventLogger, BiasExecutionContext innerBiasExecutionContext,
|
||||
ContainerExecutionContext executionContext) {
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getSubFlow()).stream()
|
||||
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle));
|
||||
String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> mainOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getSubFlow());
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> guardInputsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getGuardSubFlow()).stream()
|
||||
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, handle -> handle));
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> guardOutputHandles = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getGuardSubFlow());
|
||||
|
||||
LoopPhase phase = startPhase;
|
||||
int iterationIndex = startIterationIndex;
|
||||
Map<String, Object> inputs = startInputs;
|
||||
Map<String, Object> 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<String, Object> guardTemplateValues = buildGuardTemplateValues(inputs, latestOutputs, iterationIndex);
|
||||
Map<String, Object> 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<String, Object> 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<String, Object> 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<String, ContainerFlowInterfaceResolver.ExposedHandle> 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<ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> currentInputs, Map<String, Object> latestOutputs,
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName) {
|
||||
for (Map.Entry<String, Object> output : latestOutputs.entrySet()) {
|
||||
if (inputPortsByName.containsKey(output.getKey())) {
|
||||
currentInputs.put(output.getKey(), output.getValue());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static Map<String, Object> buildGuardTemplateValues(Map<String, Object> currentInputs,
|
||||
Map<String, Object> latestOutputs, int iteration) {
|
||||
Map<String, Object> 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<String, Object> 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<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName,
|
||||
Map<String, Object> inputValues, int iterationIndex, ContainerSubflowRole role, Map<String, Object> authorizations,
|
||||
ExecutionEventLogger eventLogger, String containerTypeName, BiasExecutionContext innerBiasExecutionContext,
|
||||
ContainerExecutionContext executionContext, Map<String, Object> continuationState,
|
||||
Map<String, Object> 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<String, Object> globalInputs = ExecutionRuntimeContextSupport.globalView(parent.getContext().getExecutionVariables());
|
||||
Map<String, ExecutionVariableDescriptor> 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<String, Object> enrichChildEventDetails(ExecutionEvent event, String innerExecutionId,
|
||||
Integer iterationIndex, String containerTypeName) {
|
||||
Map<String, Object> 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<String, Object> collectExposedOutputsAsMap(ExecutionObject execution,
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs) {
|
||||
Map<String, Object> 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<String, Object> iteratorState(List<Object> remainingValues, Map<String, Object> runtimeInputValues) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("remainingValues", new java.util.ArrayList<>(remainingValues));
|
||||
state.put("runtimeInputValues", new LinkedHashMap<>(runtimeInputValues));
|
||||
return state;
|
||||
}
|
||||
|
||||
private static Map<String, Object> loopState(Map<String, Object> currentInputs, Map<String, Object> latestOutputs) {
|
||||
Map<String, Object> state = new LinkedHashMap<>();
|
||||
state.put("currentInputs", new LinkedHashMap<>(currentInputs));
|
||||
state.put("latestOutputs", new LinkedHashMap<>(latestOutputs == null ? Map.of() : latestOutputs));
|
||||
return state;
|
||||
}
|
||||
|
||||
private static Map<String, Object> widen(Map<String, List<Object>> accumulated) {
|
||||
return new LinkedHashMap<>(accumulated);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static List<Object> asObjectList(Object value) {
|
||||
if (value instanceof List<?> list) {
|
||||
return new java.util.ArrayList<>((List<Object>) list);
|
||||
}
|
||||
return new java.util.ArrayList<>();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> asStringObjectMap(Object value) {
|
||||
if (value instanceof Map<?, ?> map) {
|
||||
return new LinkedHashMap<>((Map<String, Object>) map);
|
||||
}
|
||||
return new LinkedHashMap<>();
|
||||
}
|
||||
|
||||
private static Map<String, List<Object>> asAccumulatedOutputs(Map<String, Object> raw) {
|
||||
Map<String, List<Object>> result = new LinkedHashMap<>();
|
||||
if (raw != null) {
|
||||
raw.forEach((key, value) -> result.put(key, asObjectList(value)));
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
private record ContainerAdvanceOutcome(boolean completed, Map<String, Object> outputs, String pendingChildId) {
|
||||
static ContainerAdvanceOutcome completed(Map<String, Object> 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()
|
||||
|
|
|
|||
|
|
@ -94,7 +94,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
executionsService.prepareInput(innerExecution.getId(), handle.blockId(), handle.io().getName(), input.getValue());
|
||||
}
|
||||
|
||||
innerExecution = executionsService.startExecution(innerExecution.getId());
|
||||
innerExecution = executionsService.startContainerChild(innerExecution.getId(), executionContext);
|
||||
if (innerExecution.getContext().getStatus().isFinalState()) {
|
||||
executionsService.setContainerContinuation(executionContext.parentExecutionId(), executionContext.parentStepId(), null);
|
||||
return completedResult(innerExecution, configuration, executionVariables, executionVariableDescriptors, eventLogger);
|
||||
|
|
@ -108,7 +108,7 @@ public class GenericContainerExecutor implements ContainerExecutor<GenericContai
|
|||
.revision(2)
|
||||
.build());
|
||||
String childId = innerExecution.getId();
|
||||
return NodeExecutionResult.suspended(() -> executionsService.watchGenericSubflow(
|
||||
return NodeExecutionResult.suspended(() -> executionsService.watchContainerSubflow(
|
||||
executionContext.parentExecutionId(), executionContext.parentStepId(), childId));
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<IteratorContainerType> {
|
||||
|
||||
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<IteratorCont
|
|||
BiasExecutionContext innerBiasExecutionContext = BiasContainerPropagation.innerContextFor(container.getId(),
|
||||
List.of(configuration.getSubFlow()), biasExecutionContext);
|
||||
IteratorContainerInterfaceResolver.Resolution resolution = IteratorContainerInterfaceResolver.resolvePorts(configuration);
|
||||
IteratorContainerInterfaceResolver.ResolvedInput iteratedInput = resolution.inputsByPublicName()
|
||||
.get(resolution.iteratedPublicName());
|
||||
|
||||
Input runtimeIteratedInput = inputs.stream()
|
||||
.filter(input -> input.getDescriptor().getName().equals(resolution.iteratedPublicName()))
|
||||
.findFirst()
|
||||
|
|
@ -71,183 +55,24 @@ public class IteratorContainerExecutor implements ContainerExecutor<IteratorCont
|
|||
"IteratorContainer input " + resolution.iteratedPublicName() + " expects a list value");
|
||||
}
|
||||
|
||||
Map<String, Input> runtimeInputsByName = inputs.stream()
|
||||
.collect(LinkedHashMap::new, (map, input) -> map.put(input.getDescriptor().getName(), input), Map::putAll);
|
||||
|
||||
Map<String, List<Object>> 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<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> enrichDetails(ExecutionEvent event, String innerExecutionId, Integer iterationIndex,
|
||||
String containerType) {
|
||||
Map<String, Object> 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<IteratorContainerType> 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;
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<LoopContainerType> {
|
||||
|
||||
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<LoopContainerTyp
|
|||
BiasExecutionContext innerBiasExecutionContext = BiasContainerPropagation.innerContextFor(container.getId(),
|
||||
List.of(configuration.getSubFlow(), configuration.getGuardSubFlow()), biasExecutionContext);
|
||||
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getSubFlow()).stream()
|
||||
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed));
|
||||
String feedbackInput = resolveFeedbackInput(configuration, inputPortsByName);
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getSubFlow());
|
||||
|
||||
Map<String, Object> currentInputs = inputs.stream()
|
||||
Map<String, Object> initialInputs = inputs.stream()
|
||||
.collect(LinkedHashMap::new,
|
||||
(map, input) -> map.put(input.getDescriptor().getName(), input.getValue()),
|
||||
Map::putAll);
|
||||
Map<String, Object> 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<String, Object> 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<LoopContainerType> getContainerType() {
|
||||
return LoopContainerType.class;
|
||||
}
|
||||
|
||||
private void updateInputsForNextIteration(Map<String, Object> currentInputs, Map<String, Object> latestOutputs,
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName) {
|
||||
for (Map.Entry<String, Object> output : latestOutputs.entrySet()) {
|
||||
if (inputPortsByName.containsKey(output.getKey())) {
|
||||
currentInputs.put(output.getKey(), output.getValue());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private String resolveFeedbackInput(LoopContainerConfiguration configuration,
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> 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<ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> iterationInputs, Map<String, Object> latestOutputs, Map<String, Object> authorizations,
|
||||
Map<String, Object> executionVariables,
|
||||
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger,
|
||||
int iteration, BiasExecutionContext innerBiasExecutionContext) {
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> guardInputsByName = ContainerFlowInterfaceResolver
|
||||
.getExposedInputs(configuration.getGuardSubFlow()).stream()
|
||||
.collect(Collectors.toMap(ContainerFlowInterfaceResolver.ExposedHandle::publicName, exposed -> exposed));
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> guardOutputs = ContainerFlowInterfaceResolver
|
||||
.getExposedOutputs(configuration.getGuardSubFlow());
|
||||
Map<String, Object> guardInputValues = buildGuardTemplateValues(iterationInputs, latestOutputs, iteration);
|
||||
ExecutionObject guardExecution = executeSubFlow(
|
||||
containerName + " guard iteration " + iteration,
|
||||
configuration.getGuardSubFlow(),
|
||||
guardInputValues,
|
||||
guardInputsByName,
|
||||
"LoopContainerGuard",
|
||||
authorizations,
|
||||
executionVariables,
|
||||
executionVariableDescriptors,
|
||||
eventLogger,
|
||||
iteration,
|
||||
innerBiasExecutionContext);
|
||||
Map<String, Object> 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<String, Object> buildGuardTemplateValues(Map<String, Object> currentInputs,
|
||||
Map<String, Object> latestOutputs, int iteration) {
|
||||
Map<String, Object> 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<String, Object> inputValues,
|
||||
Map<String, ContainerFlowInterfaceResolver.ExposedHandle> inputPortsByName,
|
||||
String containerType,
|
||||
Map<String, Object> authorizations,
|
||||
Map<String, Object> executionVariables,
|
||||
Map<String, ExecutionVariableDescriptor> 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<String, Object> 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<String, ContainerFlowInterfaceResolver.ExposedHandle> 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<String, Object> collectExposedOutputs(ExecutionObject execution,
|
||||
List<ContainerFlowInterfaceResolver.ExposedHandle> exposedOutputs) {
|
||||
Map<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> enrichDetails(ExecutionEvent event, String innerExecutionId, Integer iterationIndex,
|
||||
String containerType) {
|
||||
Map<String, Object> 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) {
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -100,6 +100,17 @@ public class Step<N extends FlowNode> 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<N extends FlowNode> implements InputListener {
|
|||
private BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal();
|
||||
|
||||
@Setter
|
||||
@Getter
|
||||
@JsonIgnore
|
||||
private ExecutionEventLogger eventLogger;
|
||||
|
||||
|
|
@ -209,7 +221,8 @@ public class Step<N extends FlowNode> 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);
|
||||
|
|
|
|||
|
|
@ -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<ValidFlowStructure
|
|||
throw validationError(error(ValidationErrorCode.NESTED_CONTAINERS_NOT_SUPPORTED, "container", container.getId(), fieldPath,
|
||||
"Nested containers are not supported"));
|
||||
}
|
||||
if (!(containerConfiguration instanceof GenericContainerConfiguration)
|
||||
&& nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) {
|
||||
// GenericContainer and IteratorContainer support interactive nodes in their
|
||||
// single subflow. LoopContainer supports them only in its main subFlow, not
|
||||
// in guardSubFlow (guard interactivity is out of scope, see the interactive
|
||||
// containers plan doc).
|
||||
boolean interactiveAllowed = containerConfiguration instanceof GenericContainerConfiguration
|
||||
|| containerConfiguration instanceof IteratorContainerConfiguration
|
||||
|| (containerConfiguration instanceof LoopContainerConfiguration && "subFlow".equals(fieldName));
|
||||
if (!interactiveAllowed && nestedNodes.stream().anyMatch(FlowNode::isUserInteractive)) {
|
||||
throw validationError(error(ValidationErrorCode.INTERACTIVE_NODES_IN_CONTAINER_NOT_SUPPORTED, "container", container.getId(), fieldPath,
|
||||
"Interactive blocks inside containers are not supported yet"));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,6 +23,8 @@ import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfi
|
|||
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.flows.model.Connection;
|
||||
import it.cnr.isti.workflow.manager.flows.model.Dependency;
|
||||
import it.cnr.isti.workflow.manager.flows.model.FlowData;
|
||||
|
|
@ -125,7 +127,13 @@ public class FlowExecutionValidator {
|
|||
"specificConfiguration.subFlow",
|
||||
"Nested containers are not supported")));
|
||||
|
||||
if (!(containerConfiguration instanceof GenericContainerConfiguration)) {
|
||||
// GenericContainer and IteratorContainer support interactive nodes in
|
||||
// their single subflow; LoopContainer supports them in its main subFlow
|
||||
// (guardSubFlow interactivity is out of scope, validated separately).
|
||||
boolean interactiveAllowed = containerConfiguration instanceof GenericContainerConfiguration
|
||||
|| containerConfiguration instanceof IteratorContainerConfiguration
|
||||
|| containerConfiguration instanceof LoopContainerConfiguration;
|
||||
if (!interactiveAllowed) {
|
||||
subFlow.getNodes().stream()
|
||||
.filter(FlowNode::isUserInteractive)
|
||||
.findFirst()
|
||||
|
|
|
|||
|
|
@ -53,7 +53,6 @@ app.mcp.servers.file=${MCP_SERVERS_FILE:}
|
|||
app.executions.cache.max-size=${APP_EXECUTIONS_CACHE_MAX_SIZE:1000}
|
||||
app.executions.cache.final-ttl-ms=${APP_EXECUTIONS_CACHE_FINAL_TTL_MS:1800000}
|
||||
app.executions.cache.cleanup-interval-ms=${APP_EXECUTIONS_CACHE_CLEANUP_INTERVAL_MS:60000}
|
||||
app.container.subflow.wait-interval-ms=${CONTAINER_SUBFLOW_WAIT_INTERVAL_MS:5000}
|
||||
app.http.outbound.timeout-seconds=${HTTP_OUTBOUND_TIMEOUT_SECONDS:30}
|
||||
app.http.outbound.retry-attempts=${HTTP_OUTBOUND_RETRY_ATTEMPTS:1}
|
||||
app.http.outbound.allow-private-network=${HTTP_OUTBOUND_ALLOW_PRIVATE_NETWORK:false}
|
||||
|
|
|
|||
|
|
@ -2,7 +2,9 @@ package it.cnr.isti.workflow.manager.executions;
|
|||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
import java.util.List;
|
||||
|
|
@ -191,10 +193,21 @@ public class ExecutionTest {
|
|||
.build();
|
||||
|
||||
private FlowData loopGuardSubFlow() {
|
||||
return loopGuardSubFlow("response");
|
||||
}
|
||||
|
||||
/**
|
||||
* @param mainOutputName the exposed output name of the loop's main
|
||||
* subflow (the mocked guard prompt references it as
|
||||
* {@code outputs.<mainOutputName>}); the default
|
||||
* {@code loopGuardSubFlow()} assumes "response"
|
||||
* (the LLMBlockFactory output name).
|
||||
*/
|
||||
private FlowData loopGuardSubFlow(String mainOutputName) {
|
||||
Block<LLMBlockType> 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<SwitchBlockType> guardOutput = switchBlockFactory.create(SwitchBlockConfiguration.builder()
|
||||
.name("Expose Guard")
|
||||
|
|
@ -1412,6 +1425,209 @@ public class ExecutionTest {
|
|||
return execution;
|
||||
}
|
||||
|
||||
@Test
|
||||
public void loopContainerSuspendsForInnerInteractionAndResumesFromChild() {
|
||||
Block<HumanInteractionBlockType> innerInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("Inner review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<LoopContainerType> 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<HumanInteractionBlockType> innerInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("Inner review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<IteratorContainerType> 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<HumanInteractionBlockType> innerInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("Inner review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<GenericContainerType> 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<HumanInteractionBlockType> firstInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("First review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<GenericContainerType> 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<HumanInteractionBlockType> secondInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("Second review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<GenericContainerType> 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<HumanInteractionBlockType> innerInteraction = humanInteractiveBlockFactory.create(
|
||||
HumanInteractiveBlockConfiguration.builder()
|
||||
.name("Inner review")
|
||||
.actionDescription("Review the submitted value")
|
||||
.build());
|
||||
Container<GenericContainerType> 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<LLMBlockType> 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);
|
||||
|
|
|
|||
Loading…
Reference in New Issue